mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-03 12:41:05 +00:00
cd15ae1395
* refactor(volume): extract replica sync/select into shared volume_replica package Move the volume replica reconciliation helpers (status, union builder, SyncAndSelectBestReplica, ReadNeedleMeta) out of the shell into a new weed/storage/volume_replica package so both the shell (ec.encode, volume.tier.move, volume.check.disk) and the EC encode worker can reuse them. No behavior change. * fix(ec): bring ec.encode worker to parity with the shell - Sync replicas and encode the most-complete one (via the shared volume_replica.SyncAndSelectBestReplica) instead of a possibly-stale replica, marking all replicas readonly first. Prevents silent data loss when a stale replica is encoded and the originals deleted. - Skip remote/tiered volumes in detection (shell ec.encode excludes them). - Min-node safety gate: refuse to encode when cluster nodes < parity shards. - Align default thresholds with the shell (fullness 0.95, quiet 1h). * fix(vacuum): plugin path honors min_volume_age_seconds override deriveVacuumConfig hard-coded MinVolumeAgeSeconds=0, dropping any configured value. Read it from worker config (default 0, matching the shell/master vacuum which has no age gate) so an explicit override is honored. * address review feedback - config.go: align GetConfigSpec schema defaults (quiet_for_seconds=3600, fullness_ratio=0.95) with the runtime defaults so UI/bootstrap flows match the shell (coderabbitai). - ec_task.go: roll back readonly when markReplicasReadonly fails partway, so already-marked replicas don't stay readonly (coderabbitai). - volume_replica: pass the caller's replica statuses into buildUnionReplica instead of re-fetching them, and skip the per-needle ReadNeedleMeta RPC when the source replica is read-only (gemini-code-assist). * test(plugin_workers/ec): make fixtures eligible under the new defaults The default EC encode thresholds were raised to match the shell (fullness 0.95, quiet 1h), but the plugin-worker integration fixtures still used 90%-full / 10-minute-old volumes, so detection found no eligible volumes and the tests failed in CI. Bump the eligible fixtures to 96% full and 2h old.
109 lines
3.3 KiB
Go
109 lines
3.3 KiB
Go
package erasure_coding
|
|
|
|
import (
|
|
"context"
|
|
"io"
|
|
"net/http"
|
|
"os"
|
|
"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/stretchr/testify/require"
|
|
"google.golang.org/grpc"
|
|
"google.golang.org/grpc/credentials/insecure"
|
|
)
|
|
|
|
func TestCopyVolumeFilesToWorkerUsesCurrentCompactionRevision(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(951)
|
|
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
|
|
|
httpClient := framework.NewHTTPClient()
|
|
|
|
liveFID := framework.NewFileID(volumeID, 1001, 0x1111AAAA)
|
|
liveUploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), liveFID, []byte("live-payload-for-ec-copy"))
|
|
_ = framework.ReadAllAndClose(t, liveUploadResp)
|
|
require.Equal(t, http.StatusCreated, liveUploadResp.StatusCode)
|
|
|
|
deletedFID := framework.NewFileID(volumeID, 1002, 0x2222BBBB)
|
|
deletedUploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), deletedFID, []byte("deleted-payload-for-vacuum"))
|
|
_ = framework.ReadAllAndClose(t, deletedUploadResp)
|
|
require.Equal(t, http.StatusCreated, deletedUploadResp.StatusCode)
|
|
|
|
deleteReq, err := http.NewRequest(http.MethodDelete, clusterHarness.VolumeAdminURL()+"/"+deletedFID, nil)
|
|
require.NoError(t, err)
|
|
deleteResp := framework.DoRequest(t, httpClient, deleteReq)
|
|
_ = framework.ReadAllAndClose(t, deleteResp)
|
|
require.Equal(t, http.StatusAccepted, deleteResp.StatusCode)
|
|
|
|
compactVolumeOnce(t, grpcClient, volumeID)
|
|
|
|
task := NewErasureCodingTask(
|
|
"copy-after-compaction",
|
|
clusterHarness.VolumeServerAddress(),
|
|
volumeID,
|
|
"",
|
|
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
|
)
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer cancel()
|
|
|
|
require.NoError(t, task.markReplicasReadonly(ctx))
|
|
|
|
fileStatus, err := task.readSourceVolumeFileStatus(ctx)
|
|
require.NoError(t, err)
|
|
require.Greater(t, fileStatus.GetCompactionRevision(), uint32(0))
|
|
|
|
localFiles, err := task.copyVolumeFilesToWorker(ctx, t.TempDir())
|
|
require.NoError(t, err)
|
|
|
|
datInfo, err := os.Stat(localFiles["dat"])
|
|
require.NoError(t, err)
|
|
require.Equal(t, int64(fileStatus.GetDatFileSize()), datInfo.Size())
|
|
|
|
idxInfo, err := os.Stat(localFiles["idx"])
|
|
require.NoError(t, err)
|
|
require.Equal(t, int64(fileStatus.GetIdxFileSize()), idxInfo.Size())
|
|
}
|
|
|
|
func compactVolumeOnce(t *testing.T, grpcClient volume_server_pb.VolumeServerClient, volumeID uint32) {
|
|
t.Helper()
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer cancel()
|
|
|
|
compactStream, err := grpcClient.VacuumVolumeCompact(ctx, &volume_server_pb.VacuumVolumeCompactRequest{
|
|
VolumeId: volumeID,
|
|
})
|
|
require.NoError(t, err)
|
|
|
|
for {
|
|
_, err = compactStream.Recv()
|
|
if err == io.EOF {
|
|
break
|
|
}
|
|
require.NoError(t, err)
|
|
}
|
|
|
|
_, err = grpcClient.VacuumVolumeCommit(ctx, &volume_server_pb.VacuumVolumeCommitRequest{
|
|
VolumeId: volumeID,
|
|
})
|
|
require.NoError(t, err)
|
|
|
|
_, err = grpcClient.VacuumVolumeCleanup(ctx, &volume_server_pb.VacuumVolumeCleanupRequest{
|
|
VolumeId: volumeID,
|
|
})
|
|
require.NoError(t, err)
|
|
}
|