diff --git a/weed/shell/command_volume_check_disk.go b/weed/shell/command_volume_check_disk.go index 10e22588b..ded05c301 100644 --- a/weed/shell/command_volume_check_disk.go +++ b/weed/shell/command_volume_check_disk.go @@ -46,8 +46,10 @@ type volumeCheckDisk struct { // resurrectMissingNeedles controls whether a needle present on the source // but entirely absent on the target is pushed back. Default false: an // absent needle is indistinguishable from a vacuumed delete, so the safe - // default never raises deleted data. No caller sets this true today; it is - // the seam for a future tombstone-aware repair path. + // default never raises deleted data. Even when enabled, resurrection only + // happens into a replica whose compaction revision is 0: a never-vacuumed + // index still holds a tombstone for every delete it processed, so a needle + // absent there is provably a missing write. resurrectMissingNeedles bool ewg *ErrorWaitGroup @@ -84,6 +86,9 @@ func (c *commandVolumeCheckDisk) Help() string { -fixReadOnly: also check and repair read-only volumes using uni-directional sync -syncDeleted: sync deletion records during repair -nonRepairThreshold: maximum fraction of missing keys allowed for repair (default 0.3) + -resurrectMissingNeedles: copy needles absent on one replica back from the other, e.g. after replication failures. + Only repairs replicas that were never vacuumed (compaction revision 0), where an absent needle is provably + a missing write and not a vacuumed delete. Counts toward -nonRepairThreshold. ` } @@ -105,6 +110,7 @@ func (c *commandVolumeCheckDisk) Do(args []string, commandEnv *CommandEnv, write syncDeletions := fsckCommand.Bool("syncDeleted", false, "sync of deletions the fix") maxParallelization := fsckCommand.Int("maxParallelization", DefaultMaxParallelization, "run up to X tasks in parallel, whenever possible") nonRepairThreshold := fsckCommand.Float64("nonRepairThreshold", 0.3, "repair when missing keys is not more than this limit") + resurrectMissingNeedles := fsckCommand.Bool("resurrectMissingNeedles", false, "copy needles absent on one replica back from the other, only into never-vacuumed replicas (compaction revision 0)") if err = fsckCommand.Parse(args); err != nil { return nil } @@ -128,6 +134,8 @@ func (c *commandVolumeCheckDisk) Do(args []string, commandEnv *CommandEnv, write fixReadOnly: *fixReadOnly, nonRepairThreshold: *nonRepairThreshold, + resurrectMissingNeedles: *resurrectMissingNeedles, + ewg: NewErrorWaitGroup(*maxParallelization), } @@ -337,6 +345,21 @@ func (vcd *volumeCheckDisk) writeVerbose(format string, a ...any) { } } +// compactionRevision reads the live compaction revision of a volume replica. +// Zero means the volume has never been vacuumed. +func (vcd *volumeCheckDisk) compactionRevision(replica *VolumeReplica) (revision uint32, err error) { + err = operation.WithVolumeServerClient(false, pb.NewServerAddressFromDataNode(replica.location.dataNode), vcd.grpcDialOption(), func(client volume_server_pb.VolumeServerClient) error { + resp, reqErr := client.ReadVolumeFileStatus(context.Background(), &volume_server_pb.ReadVolumeFileStatusRequest{ + VolumeId: replica.info.Id, + }) + if resp != nil { + revision = resp.CompactionRevision + } + return reqErr + }) + return +} + // getVolumeStatusFileCount retrieves the current file count and deleted file count // from a volume server via gRPC. func (vcd *volumeCheckDisk) getVolumeStatusFileCount(vid uint32, dn *master_pb.DataNodeInfo) (totalFileCount, deletedFileCount uint64, err error) { @@ -485,9 +508,28 @@ func (vcd *volumeCheckDisk) checkBoth(source, target *VolumeReplica, bidi bool) return true, true, fmt.Errorf("readIndexDatabase %s volume %d: %w", target.location.dataNode.Id, target.info.Id, err) } + // Resurrection is gated per direction on the receiving replica's compaction + // revision: only a never-vacuumed index (revision 0) proves an absent needle + // is a missing write rather than a vacuumed delete. The revision is read + // after the index snapshot above, so the proof covers the snapshot. + resurrectIntoTarget, resurrectIntoSource := false, false + var targetRevision, sourceRevision uint32 + if vcd.resurrectMissingNeedles { + if targetRevision, err = vcd.compactionRevision(target); err != nil { + return true, true, fmt.Errorf("compactionRevision %s volume %d: %w", target.location.dataNode.Id, target.info.Id, err) + } + resurrectIntoTarget = targetRevision == 0 + if bidi { + if sourceRevision, err = vcd.compactionRevision(source); err != nil { + return true, true, fmt.Errorf("compactionRevision %s volume %d: %w", source.location.dataNode.Id, source.info.Id, err) + } + resurrectIntoSource = sourceRevision == 0 + } + } + // find and make up the differences var errs []error - targetHasChanges, errTarget := vcd.doVolumeCheckDisk(sourceDB, targetDB, source, target) + targetHasChanges, errTarget := vcd.doVolumeCheckDisk(sourceDB, targetDB, source, target, resurrectIntoTarget, targetRevision) if errTarget != nil { errs = append(errs, fmt.Errorf("doVolumeCheckDisk source:%s target:%s volume %d: %w", @@ -496,7 +538,7 @@ func (vcd *volumeCheckDisk) checkBoth(source, target *VolumeReplica, bidi bool) sourceHasChanges = false if bidi { var errSource error - sourceHasChanges, errSource = vcd.doVolumeCheckDisk(targetDB, sourceDB, target, source) + sourceHasChanges, errSource = vcd.doVolumeCheckDisk(targetDB, sourceDB, target, source, resurrectIntoSource, sourceRevision) if errSource != nil { errs = append(errs, fmt.Errorf("doVolumeCheckDisk source:%s target:%s volume %d: %w", @@ -510,7 +552,7 @@ func (vcd *volumeCheckDisk) checkBoth(source, target *VolumeReplica, bidi bool) return sourceHasChanges, targetHasChanges, nil } -func (vcd *volumeCheckDisk) doVolumeCheckDisk(minuend, subtrahend *needle_map.MemDb, source, target *VolumeReplica) (hasChanges bool, err error) { +func (vcd *volumeCheckDisk) doVolumeCheckDisk(minuend, subtrahend *needle_map.MemDb, source, target *VolumeReplica, resurrectAbsent bool, targetRevision uint32) (hasChanges bool, err error) { // find missing keys // hash join, can be more efficient @@ -533,9 +575,10 @@ func (vcd *volumeCheckDisk) doVolumeCheckDisk(minuend, subtrahend *needle_map.Me // vacuum watermark, so it cannot distinguish the two. Without // positive proof the absence is a missing write (rather than a // vacuumed delete), the safe default is to NOT resurrect: a real - // missing write may go unrepaired until a tombstone-aware path - // exists, but we never raise back data the operator deleted. - if !vcd.resurrectMissingNeedles { + // missing write may go unrepaired, but we never raise back data + // the operator deleted. resurrectAbsent carries that proof: the + // caller sets it only when the target was never vacuumed. + if !resurrectAbsent { skippedAbsentNeedles++ return nil } @@ -549,8 +592,13 @@ func (vcd *volumeCheckDisk) doVolumeCheckDisk(minuend, subtrahend *needle_map.Me }) if skippedAbsentNeedles > 0 { - vcd.write("volume %d %s: not resurrecting %d needle(s) absent on %s (cannot prove they are missing writes vs vacuumed deletes)", - source.info.Id, source.location.dataNode.Id, skippedAbsentNeedles, target.location.dataNode.Id) + if vcd.resurrectMissingNeedles { + vcd.write("volume %d %s: not resurrecting %d needle(s) absent on %s (compaction revision %d: a vacuum may have erased deleted needles there)", + source.info.Id, source.location.dataNode.Id, skippedAbsentNeedles, target.location.dataNode.Id, targetRevision) + } else { + vcd.write("volume %d %s: not resurrecting %d needle(s) absent on %s (cannot prove they are missing writes vs vacuumed deletes)", + source.info.Id, source.location.dataNode.Id, skippedAbsentNeedles, target.location.dataNode.Id) + } } vcd.write("volume %d %s has %d entries, %s missed %d and partially deleted %d entries", diff --git a/weed/shell/command_volume_check_disk_test.go b/weed/shell/command_volume_check_disk_test.go index 8a31e5cdf..2615b0b26 100644 --- a/weed/shell/command_volume_check_disk_test.go +++ b/weed/shell/command_volume_check_disk_test.go @@ -57,7 +57,7 @@ func TestDoVolumeCheckDiskDoesNotResurrectAbsentNeedle(t *testing.T) { // read, so this returns cleanly with no changes. On the pre-fix code the // needle is queued and a source blob read is attempted, which has no server // to reach and surfaces as an error (or a resurrection) instead. - hasChanges, err := vcd.doVolumeCheckDisk(sourceDB, targetDB, source, target) + hasChanges, err := vcd.doVolumeCheckDisk(sourceDB, targetDB, source, target, false, 0) if err != nil { t.Fatalf("doVolumeCheckDisk returned error: %v", err) } @@ -66,6 +66,101 @@ func TestDoVolumeCheckDiskDoesNotResurrectAbsentNeedle(t *testing.T) { } } +// TestDoVolumeCheckDiskResurrectSkipsVacuumedTarget verifies that even with +// -resurrectMissingNeedles, an absent needle is not pushed into a replica that +// has been vacuumed (compaction revision > 0): the vacuum may have erased the +// tombstone of a legitimate delete, so the absence proves nothing. +func TestDoVolumeCheckDiskResurrectSkipsVacuumedTarget(t *testing.T) { + sourceDB, targetDB := needle_map.NewMemDb(), needle_map.NewMemDb() + defer sourceDB.Close() + defer targetDB.Close() + + if err := sourceDB.Set(types.NeedleId(1001), types.ToOffset(8), types.Size(123)); err != nil { + t.Fatalf("seed source: %v", err) + } + + var buf bytes.Buffer + vcd := &volumeCheckDisk{ + commandEnv: &CommandEnv{ + option: &ShellOptions{ + GrpcDialOption: grpc.WithTransportCredentials(insecure.NewCredentials()), + }, + }, + writer: &buf, + now: time.Now(), + applyChanges: false, + nonRepairThreshold: 1, + resurrectMissingNeedles: true, + } + + source := &VolumeReplica{ + location: &location{"dc1", "r1", &master_pb.DataNodeInfo{Id: "127.0.0.1:1"}}, + info: &master_pb.VolumeInformationMessage{Id: 7}, + } + target := &VolumeReplica{ + location: &location{"dc1", "r2", &master_pb.DataNodeInfo{Id: "127.0.0.1:2"}}, + info: &master_pb.VolumeInformationMessage{Id: 7}, + } + + // resurrectAbsent=false with revision 2: the caller found the target vacuumed. + hasChanges, err := vcd.doVolumeCheckDisk(sourceDB, targetDB, source, target, false, 2) + if err != nil { + t.Fatalf("doVolumeCheckDisk returned error: %v", err) + } + if hasChanges { + t.Fatalf("absent needle was resurrected into a vacuumed replica; expected no changes") + } + if !bytes.Contains(buf.Bytes(), []byte("compaction revision 2")) { + t.Fatalf("expected skip message to mention the compaction revision, got: %s", buf.String()) + } +} + +// TestDoVolumeCheckDiskResurrectQueuesProvenMissingWrite verifies that with +// -resurrectMissingNeedles and a never-vacuumed target, the absent needle is +// treated as a missing write and queued for repair: the blob read against the +// unreachable source server is attempted, which is only possible if the needle +// passed the resurrection guard. +func TestDoVolumeCheckDiskResurrectQueuesProvenMissingWrite(t *testing.T) { + sourceDB, targetDB := needle_map.NewMemDb(), needle_map.NewMemDb() + defer sourceDB.Close() + defer targetDB.Close() + + if err := sourceDB.Set(types.NeedleId(1001), types.ToOffset(8), types.Size(123)); err != nil { + t.Fatalf("seed source: %v", err) + } + + var buf bytes.Buffer + vcd := &volumeCheckDisk{ + commandEnv: &CommandEnv{ + option: &ShellOptions{ + GrpcDialOption: grpc.WithTransportCredentials(insecure.NewCredentials()), + }, + }, + writer: &buf, + now: time.Now(), + applyChanges: false, + nonRepairThreshold: 1, + resurrectMissingNeedles: true, + } + + source := &VolumeReplica{ + location: &location{"dc1", "r1", &master_pb.DataNodeInfo{Id: "127.0.0.1:1"}}, + info: &master_pb.VolumeInformationMessage{Id: 7}, + } + target := &VolumeReplica{ + location: &location{"dc1", "r2", &master_pb.DataNodeInfo{Id: "127.0.0.1:2"}}, + info: &master_pb.VolumeInformationMessage{Id: 7}, + } + + _, err := vcd.doVolumeCheckDisk(sourceDB, targetDB, source, target, true, 0) + if err == nil { + t.Fatalf("expected blob read against unreachable source, got nil error (needle did not pass the guard?)") + } + if !bytes.Contains(buf.Bytes(), []byte("missed 1")) { + t.Fatalf("expected the absent needle to be counted as missed, got: %s", buf.String()) + } +} + type testCommandVolumeCheckDisk struct { commandVolumeCheckDisk } diff --git a/weed/shell/command_volume_fix_replication.go b/weed/shell/command_volume_fix_replication.go index 526256e5f..c2277d69d 100644 --- a/weed/shell/command_volume_fix_replication.go +++ b/weed/shell/command_volume_fix_replication.go @@ -260,7 +260,7 @@ func checkOneVolume(a *VolumeReplica, b *VolumeReplica, writer io.Writer, comman if err := vcd.readIndexDatabase(bDB, b.info.Collection, b.info.Id, pb.NewServerAddressFromDataNode(b.location.dataNode)); err != nil { return fmt.Errorf("readIndexDatabase %s volume %d: %v", b.location.dataNode, b.info.Id, err) } - if _, err = vcd.doVolumeCheckDisk(aDB, bDB, a, b); err != nil { + if _, err = vcd.doVolumeCheckDisk(aDB, bDB, a, b, false, 0); err != nil { return fmt.Errorf("doVolumeCheckDisk source:%s target:%s volume %d: %v", a.location.dataNode.Id, b.location.dataNode.Id, a.info.Id, err) } return