diff --git a/weed/worker/tasks/erasure_coding/ec_task.go b/weed/worker/tasks/erasure_coding/ec_task.go index 7bdc4f741..0a49e3ac1 100644 --- a/weed/worker/tasks/erasure_coding/ec_task.go +++ b/weed/worker/tasks/erasure_coding/ec_task.go @@ -596,7 +596,6 @@ func (t *ErasureCodingTask) deleteOriginalVolume(ctx context.Context) error { } } - // Report results if len(deleteErrors) > 0 { t.GetLogger().WithFields(map[string]interface{}{ "volume_id": t.volumeID, @@ -605,16 +604,19 @@ func (t *ErasureCodingTask) deleteOriginalVolume(ctx context.Context) error { "total_replicas": len(replicas), "success_rate": float64(successCount) / float64(len(replicas)) * 100, "errors": deleteErrors, - }).Warning("Some volume deletions failed") - // Don't return error - EC task should still be considered successful if shards are mounted - } else { - t.GetLogger().WithFields(map[string]interface{}{ - "volume_id": t.volumeID, - "replica_count": len(replicas), - "replica_servers": replicas, - }).Info("Successfully deleted volume from all replica servers") + }).Error("Failed to delete some original volume replicas after EC encoding") + // A surviving source replica lets a later detection scan re-propose + // EC on the same volume, which retries over mounted shards. + return fmt.Errorf("failed to delete %d of %d original volume replicas for volume %d: %s", + len(deleteErrors), len(replicas), t.volumeID, strings.Join(deleteErrors, "; ")) } + t.GetLogger().WithFields(map[string]interface{}{ + "volume_id": t.volumeID, + "replica_count": len(replicas), + "replica_servers": replicas, + }).Info("Successfully deleted volume from all replica servers") + return nil } diff --git a/weed/worker/tasks/erasure_coding/ec_task_delete_swallow_test.go b/weed/worker/tasks/erasure_coding/ec_task_delete_swallow_test.go new file mode 100644 index 000000000..54c5f317a --- /dev/null +++ b/weed/worker/tasks/erasure_coding/ec_task_delete_swallow_test.go @@ -0,0 +1,113 @@ +package erasure_coding + +import ( + "context" + "net/http" + "strings" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/test/volume_server/framework" + "github.com/seaweedfs/seaweedfs/test/volume_server/matrix" + "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" + "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" + "github.com/stretchr/testify/require" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" +) + +// One reachable replica + one unreachable: the reachable delete still +// succeeds, and the function surfaces an error naming the failure. +func TestDeleteOriginalVolumeSurfacesReplicaFailures(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test in short mode") + } + + clusterHarness := framework.StartVolumeCluster(t, matrix.P1()) + conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress()) + defer conn.Close() + + const volumeID = uint32(91842) + framework.AllocateVolume(t, grpcClient, volumeID, "") + + httpClient := framework.NewHTTPClient() + fid := framework.NewFileID(volumeID, 918420, 0x91842042) + uploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), fid, + []byte("delete-surface-content-for-issue-9184")) + _ = framework.ReadAllAndClose(t, uploadResp) + require.Equal(t, http.StatusCreated, uploadResp.StatusCode) + + task := NewErasureCodingTask( + "delete-surface-fix", + clusterHarness.VolumeServerAddress(), + volumeID, + "", + grpc.WithTransportCredentials(insecure.NewCredentials()), + ) + + unreachable := "127.0.0.1:1" + task.sources = []*worker_pb.TaskSource{ + { + Node: clusterHarness.VolumeServerAddress(), + VolumeId: volumeID, + }, + { + Node: unreachable, + VolumeId: volumeID, + }, + } + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + err := task.deleteOriginalVolume(ctx) + require.Error(t, err, "deleteOriginalVolume must surface replica delete failures (#9184)") + require.Contains(t, err.Error(), unreachable, + "returned error should name the replica that failed: %v", err) + require.True(t, + strings.Contains(err.Error(), "failed to delete"), + "returned error should describe what failed: %v", err) + + _, statusErr := grpcClient.VolumeStatus(ctx, &volume_server_pb.VolumeStatusRequest{VolumeId: volumeID}) + require.Error(t, statusErr, + "reachable replica %d should have been deleted before failure was surfaced", volumeID) +} + +func TestDeleteOriginalVolumeSucceedsWhenAllReplicasReachable(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test in short mode") + } + + clusterHarness := framework.StartVolumeCluster(t, matrix.P1()) + conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress()) + defer conn.Close() + + const volumeID = uint32(91844) + framework.AllocateVolume(t, grpcClient, volumeID, "") + + httpClient := framework.NewHTTPClient() + fid := framework.NewFileID(volumeID, 918440, 0x91844042) + uploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), fid, + []byte("delete-happy-path-content-for-issue-9184")) + _ = framework.ReadAllAndClose(t, uploadResp) + require.Equal(t, http.StatusCreated, uploadResp.StatusCode) + + task := NewErasureCodingTask( + "delete-happy-path", + clusterHarness.VolumeServerAddress(), + volumeID, + "", + grpc.WithTransportCredentials(insecure.NewCredentials()), + ) + task.sources = []*worker_pb.TaskSource{ + {Node: clusterHarness.VolumeServerAddress(), VolumeId: volumeID}, + } + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + require.NoError(t, task.deleteOriginalVolume(ctx)) + + _, statusErr := grpcClient.VolumeStatus(ctx, &volume_server_pb.VolumeStatusRequest{VolumeId: volumeID}) + require.Error(t, statusErr, "volume %d should be gone after successful delete", volumeID) +}