diff --git a/seaweed-volume/proto/volume_server.proto b/seaweed-volume/proto/volume_server.proto index f23f05863..5ae5511e4 100644 --- a/seaweed-volume/proto/volume_server.proto +++ b/seaweed-volume/proto/volume_server.proto @@ -168,6 +168,7 @@ message VacuumVolumeCheckRequest { } message VacuumVolumeCheckResponse { double garbage_ratio = 1; + bool disk_space_low = 4; // the volume is read-only solely because its disk is low on space — a cause compaction itself reclaims } message VacuumVolumeCompactRequest { diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 3c59d8e38..7d36c0443 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -1127,14 +1127,24 @@ impl VolumeServer for VolumeGrpcService { ) -> Result, Status> { let vid = VolumeId(request.into_inner().volume_id); let store = self.state.store.read().unwrap(); - let garbage_ratio = match store.find_volume(vid) { - Some((_, vol)) => vol.garbage_level(), + let (garbage_ratio, disk_space_low) = match store.find_volume(vid) { + Some((_, vol)) => { + // disk_space_low only counts when it is the sole read-only + // cause — an operator mark or I/O quarantine still shields + // the volume. + let (_, no_write_or_delete, no_write_can_delete, is_low) = vol.read_only_reasons(); + ( + vol.garbage_level(), + is_low && !no_write_or_delete && !no_write_can_delete, + ) + } None => { return Err(crate::storage::volume::VolumeError::VolumeNotFound(vid).into()); } }; Ok(Response::new(volume_server_pb::VacuumVolumeCheckResponse { garbage_ratio, + disk_space_low, })) } diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index 0c656888a..29d254aa9 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -14,9 +14,9 @@ use crate::pb::master_pb; use crate::storage::disk_location::DiskLocation; use crate::storage::erasure_coding::ec_shard::{EcVolumeShard, MAX_SHARD_COUNT, ShardId}; use crate::storage::erasure_coding::ec_volume::{EcVolume, is_usable_ecx_file}; -use crate::storage::needle::needle::Needle; +use crate::storage::needle::needle::{Needle, get_actual_size}; use crate::storage::needle_map::NeedleMapKind; -use crate::storage::super_block::ReplicaPlacement; +use crate::storage::super_block::{ReplicaPlacement, SUPER_BLOCK_SIZE}; use crate::storage::types::*; use crate::storage::volume::{CompactionJob, VifVolumeInfo, VolumeError, VolumeSpec}; @@ -1539,13 +1539,21 @@ impl Store { preallocate: u64, ) -> Result, VolumeError> { // Required space matches Go's CompactVolume check: the larger of the - // requested preallocation and the estimated volume size. + // requested preallocation and the estimated compacted size — the live + // needles, not the .dat the garbage already occupies, so a full disk + // can still be reclaimed. let (loc_idx, space_needed) = { let (loc_idx, v) = self .find_volume(vid) .ok_or(VolumeError::VolumeNotFound(vid))?; - let estimated = v.dat_file_size().unwrap_or(0) + v.idx_file_size(); - (loc_idx, std::cmp::max(preallocate, estimated)) + let live_count = (v.file_count() - v.deleted_count()).max(0) as u64; + let live_bytes = v.content_size().saturating_sub(v.deleted_size()); + let per_needle = (get_actual_size(Size(0), v.version()) + + NEEDLE_PADDING_SIZE as i64 + + NEEDLE_MAP_ENTRY_SIZE as i64) as u64; + let estimated = SUPER_BLOCK_SIZE as u64 + live_count * per_needle + live_bytes; + let space_needed = std::cmp::max(preallocate, estimated); + (loc_idx, space_needed + space_needed / 10) }; let dir = self.locations[loc_idx].directory.clone(); diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index d36261d71..4e6ee064e 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -2690,6 +2690,19 @@ impl Volume { || self.location_disk_space_low.load(Ordering::Relaxed) } + /// Mirrors Go's ReadOnlyReasons: `no_write_or_delete` already covers the + /// io_unavailable quarantine. + pub fn read_only_reasons(&self) -> (bool, bool, bool, bool) { + let no_write_or_delete = self.no_write_or_delete || self.io_unavailable.is_some(); + let disk_space_low = self.location_disk_space_low.load(Ordering::Relaxed); + ( + no_write_or_delete || self.no_write_can_delete || disk_space_low, + no_write_or_delete, + self.no_write_can_delete, + disk_space_low, + ) + } + /// The reason the volume refuses all I/O, when a failed recovery left the /// .dat/index pair unverified. Mirrors Go's unavailableError. pub fn unavailable_error(&self) -> Option { diff --git a/weed/pb/volume_server.proto b/weed/pb/volume_server.proto index abbf6de00..7a20d5632 100644 --- a/weed/pb/volume_server.proto +++ b/weed/pb/volume_server.proto @@ -168,6 +168,7 @@ message VacuumVolumeCheckRequest { } message VacuumVolumeCheckResponse { double garbage_ratio = 1; + bool disk_space_low = 4; // the volume is read-only solely because its disk is low on space — a cause compaction itself reclaims } message VacuumVolumeCompactRequest { diff --git a/weed/pb/volume_server_pb/volume_server.pb.go b/weed/pb/volume_server_pb/volume_server.pb.go index 0c568858d..653bea272 100644 --- a/weed/pb/volume_server_pb/volume_server.pb.go +++ b/weed/pb/volume_server_pb/volume_server.pb.go @@ -436,6 +436,7 @@ func (x *VacuumVolumeCheckRequest) GetVolumeId() uint32 { type VacuumVolumeCheckResponse struct { state protoimpl.MessageState `protogen:"open.v1"` GarbageRatio float64 `protobuf:"fixed64,1,opt,name=garbage_ratio,json=garbageRatio,proto3" json:"garbage_ratio,omitempty"` + DiskSpaceLow bool `protobuf:"varint,4,opt,name=disk_space_low,json=diskSpaceLow,proto3" json:"disk_space_low,omitempty"` // the volume is read-only solely because its disk is low on space — a cause compaction itself reclaims unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -477,6 +478,13 @@ func (x *VacuumVolumeCheckResponse) GetGarbageRatio() float64 { return 0 } +func (x *VacuumVolumeCheckResponse) GetDiskSpaceLow() bool { + if x != nil { + return x.DiskSpaceLow + } + return false +} + type VacuumVolumeCompactRequest struct { state protoimpl.MessageState `protogen:"open.v1"` VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"` @@ -7276,9 +7284,10 @@ const file_volume_server_proto_rawDesc = "" + "\aversion\x18\x05 \x01(\rR\aversion\"\a\n" + "\x05Empty\"7\n" + "\x18VacuumVolumeCheckRequest\x12\x1b\n" + - "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\"@\n" + + "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\"f\n" + "\x19VacuumVolumeCheckResponse\x12#\n" + - "\rgarbage_ratio\x18\x01 \x01(\x01R\fgarbageRatio\"[\n" + + "\rgarbage_ratio\x18\x01 \x01(\x01R\fgarbageRatio\x12$\n" + + "\x0edisk_space_low\x18\x04 \x01(\bR\fdiskSpaceLow\"[\n" + "\x1aVacuumVolumeCompactRequest\x12\x1b\n" + "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12 \n" + "\vpreallocate\x18\x02 \x01(\x03R\vpreallocate\"f\n" + diff --git a/weed/server/volume_grpc_vacuum.go b/weed/server/volume_grpc_vacuum.go index 7bea72497..e56c48405 100644 --- a/weed/server/volume_grpc_vacuum.go +++ b/weed/server/volume_grpc_vacuum.go @@ -22,9 +22,10 @@ func (vs *VolumeServer) VacuumVolumeCheck(ctx context.Context, req *volume_serve resp := &volume_server_pb.VacuumVolumeCheckResponse{} - garbageRatio, err := vs.store.CheckCompactVolume(needle.VolumeId(req.VolumeId)) + garbageRatio, diskSpaceLow, err := vs.store.CheckCompactVolume(needle.VolumeId(req.VolumeId)) resp.GarbageRatio = garbageRatio + resp.DiskSpaceLow = diskSpaceLow if err != nil { glog.V(3).Infof("check volume %d: %v", req.VolumeId, err) diff --git a/weed/storage/store_vacuum.go b/weed/storage/store_vacuum.go index 8eb69df62..99878c952 100644 --- a/weed/storage/store_vacuum.go +++ b/weed/storage/store_vacuum.go @@ -7,16 +7,21 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" + "github.com/seaweedfs/seaweedfs/weed/storage/types" ) var ErrInsufficientSpace = fmt.Errorf("insufficient free space") -func (s *Store) CheckCompactVolume(volumeId needle.VolumeId) (float64, error) { +func (s *Store) CheckCompactVolume(volumeId needle.VolumeId) (garbageRatio float64, diskSpaceLow bool, err error) { if v := s.findVolume(volumeId); v != nil { glog.V(3).Infof("volume %d garbage level: %f", volumeId, v.garbageLevel()) - return v.garbageLevel(), nil + // diskSpaceLow only counts when it is the sole read-only cause — an + // operator mark or I/O quarantine still shields the volume. + _, noWriteOrDelete, noWriteCanDelete, isLow := v.ReadOnlyReasons() + return v.garbageLevel(), isLow && !noWriteOrDelete && !noWriteCanDelete, nil } - return 0, fmt.Errorf("volume id %d is not found during check compact: %w", volumeId, ErrVolumeNotFound) + return 0, false, fmt.Errorf("volume id %d is not found during check compact: %w", volumeId, ErrVolumeNotFound) } func (s *Store) CompactVolume(vid needle.VolumeId, preallocate int64, compactionBytePerSecond int64, progressFn ProgressFunc) error { @@ -56,18 +61,39 @@ func (s *Store) CommitCleanupVolume(vid needle.VolumeId) error { return fmt.Errorf("volume id %d is not found during cleaning up: %w", vid, ErrVolumeNotFound) } +// estimatedCompactedSize is what compaction writes: a superblock, the live +// needles with their on-disk framing, and an index with live entries only. +// Deleted bytes do not carry over, so a mostly-garbage volume needs far less +// space than it occupies. +func estimatedCompactedSize(v *Volume) int64 { + liveCount := v.FileCount() + if deleted := v.DeletedCount(); deleted < liveCount { + liveCount -= deleted + } else { + liveCount = 0 + } + liveBytes := v.ContentSize() + if deleted := v.DeletedSize(); deleted < liveBytes { + liveBytes -= deleted + } else { + liveBytes = 0 + } + perNeedle := needle.GetActualSize(0, v.Version()) + types.NeedlePaddingSize + types.NeedleMapEntrySize + return super_block.SuperBlockSize + int64(liveCount)*perNeedle + int64(liveBytes) +} + func ensureCompactVolumeSpace(v *Volume, preallocate int64) error { - // Get current volume size for space calculation volumeSize, indexSize, _ := v.FileStat() - // Calculate space needed for compaction: - // 1. Space for the new compacted volume (approximately same as current volume size) - // 2. Use the larger of preallocate or estimated volume size - estimatedCompactSize := int64(volumeSize + indexSize) + // The compacted output holds live needles only, so measure against the + // estimated compacted size — otherwise a disk full of garbage can never + // reclaim itself. + estimatedCompactSize := estimatedCompactedSize(v) spaceNeeded := preallocate if estimatedCompactSize > preallocate { spaceNeeded = estimatedCompactSize } + spaceNeeded += spaceNeeded / 10 diskStatus := stats.NewDiskStatus(v.dir) if int64(diskStatus.Free) < spaceNeeded { diff --git a/weed/storage/store_vacuum_test.go b/weed/storage/store_vacuum_test.go index ae3293461..8d7bb7139 100644 --- a/weed/storage/store_vacuum_test.go +++ b/weed/storage/store_vacuum_test.go @@ -2,6 +2,12 @@ package storage import ( "testing" + + "github.com/stretchr/testify/require" + + "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" + "github.com/seaweedfs/seaweedfs/weed/storage/types" ) func TestSpaceCalculation(t *testing.T) { @@ -49,3 +55,104 @@ func TestSpaceCalculation(t *testing.T) { }) } } + +// Compaction writes live needles only, so the space check must be measured +// against the live size, not the .dat the garbage occupies — a full disk +// needs the estimate to shrink or it can never reclaim. +func TestEstimatedCompactedSizeCountsLiveNeedles(t *testing.T) { + dir := t.TempDir() + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + defer v.Close() + + const count = 20 + for i := 1; i <= count; i++ { + if _, _, _, err := v.writeNeedle2(newRandomNeedle(uint64(i)), true, false, false); err != nil { + t.Fatalf("write needle %d: %v", i, err) + } + } + datSize, _, _ := v.FileStat() + + fullEstimate := estimatedCompactedSize(v) + if fullEstimate <= super_block.SuperBlockSize { + t.Fatalf("estimate for all-live volume = %d, want > superblock", fullEstimate) + } + + for i := 1; i < count; i++ { + if _, err := v.doDeleteRequest(newEmptyNeedle(uint64(i))); err != nil { + t.Fatalf("delete needle %d: %v", i, err) + } + } + + estimate := estimatedCompactedSize(v) + if estimate >= int64(datSize) { + t.Fatalf("estimate %d not below .dat size %d with 19/20 needles deleted", estimate, datSize) + } + if estimate <= super_block.SuperBlockSize { + t.Fatalf("estimate %d lost the one live needle", estimate) + } + live := int64(v.FileCount()-v.DeletedCount())*types.NeedleMapEntrySize + super_block.SuperBlockSize + if estimate < live { + t.Fatalf("estimate %d below superblock + live index entries %d", estimate, live) + } + + if _, err := v.doDeleteRequest(newEmptyNeedle(uint64(count))); err != nil { + t.Fatalf("delete last needle: %v", err) + } + if estimate := estimatedCompactedSize(v); estimate != super_block.SuperBlockSize { + t.Fatalf("all-deleted estimate = %d, want superblock only (%d)", estimate, super_block.SuperBlockSize) + } +} + +// The estimate must cover what compaction writes on disk: each live needle's +// content plus its header, checksum, timestamp and padding. An all-live +// volume's compacted .dat is byte-for-byte its current one, so the estimate +// may not fall below the current file. +func TestEstimatedCompactedSizeCoversNeedleFraming(t *testing.T) { + dir := t.TempDir() + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + defer v.Close() + + for i := 1; i <= 100; i++ { + if _, _, _, err := v.writeNeedle2(newRandomNeedle(uint64(i)), true, false, false); err != nil { + t.Fatalf("write needle %d: %v", i, err) + } + } + datSize, _, _ := v.FileStat() + + if estimate := estimatedCompactedSize(v); estimate < int64(datSize) { + t.Fatalf("estimate %d below .dat size %d for an all-live volume: missing per-needle framing", estimate, datSize) + } +} + +// disk_space_low is only reported when low space is the sole read-only cause, +// so a volume also marked read-only by an operator or quarantined by failed +// I/O stays out of the sweep. +func TestCheckCompactVolumeDiskLowSoleCauseOnly(t *testing.T) { + dir := t.TempDir() + store := newSingleDirStore(t, dir) + defer store.Close() + const vid = needle.VolumeId(7) + require.NoError(t, store.AddVolume(vid, "", NeedleMapInMemory, "000", "", 0, needle.GetCurrentVersion(), 0, types.HardDriveType, 0)) + + _, low, err := store.CheckCompactVolume(vid) + require.NoError(t, err) + require.False(t, low) + + store.Locations[0].isDiskSpaceLow.Store(true) + _, low, err = store.CheckCompactVolume(vid) + require.NoError(t, err) + require.True(t, low) + + require.NoError(t, store.MarkVolumeReadonly(vid, false, false)) + _, low, err = store.CheckCompactVolume(vid) + require.NoError(t, err) + require.False(t, low) +} diff --git a/weed/topology/topology_vacuum.go b/weed/topology/topology_vacuum.go index 0a4537a94..4480637d2 100644 --- a/weed/topology/topology_vacuum.go +++ b/weed/topology/topology_vacuum.go @@ -23,11 +23,11 @@ import ( ) func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid needle.VolumeId, - locationlist *VolumeLocationList, garbageThreshold float64) (*VolumeLocationList, bool) { + locationlist *VolumeLocationList, garbageThreshold float64, skipReadOnly bool) (*VolumeLocationList, bool) { ch := make(chan int, locationlist.Length()) errCount := int32(0) for index, dn := range locationlist.list { - go func(index int, url pb.ServerAddress, vid needle.VolumeId) { + go func(index int, dn *DataNode, url pb.ServerAddress, vid needle.VolumeId) { err := operation.WithVolumeServerClient(false, url, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error { resp, err := volumeServerClient.VacuumVolumeCheck(context.Background(), &volume_server_pb.VacuumVolumeCheckRequest{ VolumeId: uint32(vid), @@ -37,6 +37,21 @@ func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid ne ch <- -1 return err } + // A sweep skips a read-only copy unless low disk space is its + // only read-only cause — that is the copy compaction exists for. + if skipReadOnly { + v, lookErr := dn.GetVolumesById(vid) + if lookErr != nil { + atomic.AddInt32(&errCount, 1) + ch <- -1 + return lookErr + } + if v.ReadOnly && !resp.DiskSpaceLow { + glog.V(0).Infof("skip vacuuming read-only volume %d on %s", vid, url) + ch <- -1 + return nil + } + } if resp.GarbageRatio >= garbageThreshold { ch <- index } else { @@ -47,7 +62,7 @@ func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid ne if err != nil { glog.V(0).Infof("Checking vacuuming %d on %s: %v", vid, url, err) } - }(index, dn.ServerAddress(), vid) + }(index, dn, dn.ServerAddress(), vid) } vacuumLocationList := NewVolumeLocationList() @@ -356,18 +371,16 @@ func (t *Topology) vacuumOneVolumeLayout(grpcDialOption grpc.DialOption, volumeL } // skipReadOnly is set by the background scan and all-volumes sweep, where a -// read-only flag usually means an unhealthy disk. An explicit volumeId clears -// it so a benignly read-only (full/oversized) volume can be reclaimed. +// read-only flag usually means an unhealthy disk. Even then a copy that is +// read-only because its disk is low on space stays eligible — compaction is +// how the space comes back. An explicit volumeId clears the rule entirely. func (t *Topology) vacuumOneVolumeId(grpcDialOption grpc.DialOption, volumeLayout *VolumeLayout, c *Collection, garbageThreshold float64, locationList *VolumeLocationList, vid needle.VolumeId, preallocate int64, skipReadOnly bool) { volumeLayout.accessLock.RLock() isReadOnly := volumeLayout.vid2location[vid].AnyReadOnly() isEnoughCopies := volumeLayout.enoughCopies(vid) volumeLayout.accessLock.RUnlock() - if isReadOnly { - if skipReadOnly { - return - } + if isReadOnly && !skipReadOnly { glog.V(0).Infof("vacuuming read-only volume %d on explicit request", vid) } if !isEnoughCopies { @@ -377,7 +390,7 @@ func (t *Topology) vacuumOneVolumeId(grpcDialOption grpc.DialOption, volumeLayou glog.V(1).Infof("check vacuum on collection:%s volume:%d", c.Name, vid) if vacuumLocationList, needVacuum := t.batchVacuumVolumeCheck( - grpcDialOption, vid, locationList, garbageThreshold); needVacuum { + grpcDialOption, vid, locationList, garbageThreshold, skipReadOnly); needVacuum { if t.batchVacuumVolumeCompact(grpcDialOption, volumeLayout, vid, vacuumLocationList, preallocate) { t.batchVacuumVolumeCommit(grpcDialOption, volumeLayout, vid, vacuumLocationList, locationList) } else { diff --git a/weed/topology/topology_vacuum_test.go b/weed/topology/topology_vacuum_test.go index fd193b495..2b504b516 100644 --- a/weed/topology/topology_vacuum_test.go +++ b/weed/topology/topology_vacuum_test.go @@ -64,7 +64,7 @@ func (f *fakeVolumeDeleteServer) VolumeDelete(ctx context.Context, req *volume_s return &volume_server_pb.VolumeDeleteResponse{}, nil } -func startFakeVolumeServer(t *testing.T, vs *fakeVolumeDeleteServer) (grpcPort int, dialOption grpc.DialOption) { +func startFakeVolumeServer(t *testing.T, vs volume_server_pb.VolumeServerServer) (grpcPort int, dialOption grpc.DialOption) { t.Helper() lis, err := net.Listen("tcp", "127.0.0.1:0") if err != nil { @@ -253,3 +253,110 @@ func TestDeleteEmptyVolumesKeepsVidWhenCopyDeleteFails(t *testing.T) { t.Fatalf("VolumeDelete calls = %d, want 1 (only the reachable copy)", len(fake.deletes)) } } + +type fakeVacuumServer struct { + volume_server_pb.UnimplementedVolumeServerServer + mu sync.Mutex + checks map[uint32]*volume_server_pb.VacuumVolumeCheckResponse + committed []uint32 +} + +func (f *fakeVacuumServer) VacuumVolumeCheck(_ context.Context, req *volume_server_pb.VacuumVolumeCheckRequest) (*volume_server_pb.VacuumVolumeCheckResponse, error) { + resp, ok := f.checks[req.VolumeId] + if !ok { + return nil, fmt.Errorf("volume %d not found", req.VolumeId) + } + return resp, nil +} + +func (f *fakeVacuumServer) VacuumVolumeCompact(_ *volume_server_pb.VacuumVolumeCompactRequest, _ volume_server_pb.VolumeServer_VacuumVolumeCompactServer) error { + return nil +} + +func (f *fakeVacuumServer) VacuumVolumeCommit(_ context.Context, req *volume_server_pb.VacuumVolumeCommitRequest) (*volume_server_pb.VacuumVolumeCommitResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.committed = append(f.committed, req.VolumeId) + return &volume_server_pb.VacuumVolumeCommitResponse{IsReadOnly: true}, nil +} + +func (f *fakeVacuumServer) VacuumVolumeCleanup(_ context.Context, _ *volume_server_pb.VacuumVolumeCleanupRequest) (*volume_server_pb.VacuumVolumeCleanupResponse, error) { + return &volume_server_pb.VacuumVolumeCleanupResponse{}, nil +} + +// A sweep keeps a read-only replica eligible only when the volume server +// reports disk_space_low; other read-only causes stay skipped unless the +// request names the volume explicitly. +func TestVacuumReadOnlyDiskLowVolume(t *testing.T) { + fake := &fakeVacuumServer{ + checks: map[uint32]*volume_server_pb.VacuumVolumeCheckResponse{ + 1: {GarbageRatio: 0.9, DiskSpaceLow: true}, + 2: {GarbageRatio: 0.9}, + 3: {GarbageRatio: 0.9}, + 4: {GarbageRatio: 0.1, DiskSpaceLow: true}, + 5: {GarbageRatio: 0.9}, + }, + } + grpcPort, dialOption := startFakeVolumeServer(t, fake) + + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + dn := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1"). + GetOrCreateDataNode("127.0.0.1", 8080, grpcPort, "127.0.0.1", "dn", map[string]uint32{"": 10}) + + newVolumeInfo := func(vid needle.VolumeId, readOnly bool) storage.VolumeInfo { + return storage.VolumeInfo{ + Id: vid, + Size: 1 << 20, + Collection: "c", + ReadOnly: readOnly, + ModifiedAtSecond: time.Now().Unix(), + Version: needle.GetCurrentVersion(), + ReplicaPlacement: &super_block.ReplicaPlacement{}, + Ttl: needle.EMPTY_TTL, + } + } + vl := topo.GetVolumeLayout("c", &super_block.ReplicaPlacement{}, needle.EMPTY_TTL, types.ToDiskType("")) + volumes := []storage.VolumeInfo{ + newVolumeInfo(1, true), + newVolumeInfo(2, true), + newVolumeInfo(3, false), + newVolumeInfo(4, true), + newVolumeInfo(5, true), + } + dn.UpdateVolumes(volumes) + for _, v := range volumes { + topo.RegisterVolumeLayout(v, dn) + } + c := NewCollection("c", topo.volumeSizeLimit, false) + + vacuum := func(vid needle.VolumeId, skipReadOnly bool) { + vl.accessLock.RLock() + ll := vl.vid2location[vid].Copy() + vl.accessLock.RUnlock() + topo.vacuumOneVolumeId(dialOption, vl, c, 0.3, ll, vid, 0, skipReadOnly) + } + + vacuum(1, true) // read-only but disk_low: compacted + vacuum(2, true) // read-only otherwise: skipped by sweep + vacuum(3, true) // writable: compacted + vacuum(4, true) // disk_low but below threshold: skipped + vacuum(5, false) // read-only, explicit request: compacted + + fake.mu.Lock() + defer fake.mu.Unlock() + want := map[uint32]bool{1: true, 3: true, 5: true} + got := map[uint32]bool{} + for _, vid := range fake.committed { + got[vid] = true + } + for vid := range want { + if !got[vid] { + t.Errorf("volume %d was not vacuum-committed", vid) + } + } + for _, vid := range fake.committed { + if !want[vid] { + t.Errorf("volume %d vacuum-committed unexpectedly", vid) + } + } +}