fix(ec): wipe stale EC artifacts before re-encoding

VolumeEcShardsGenerate evicts any in-memory EcVolume for the volume and
removes its on-disk shard/index/sidecar files before writing fresh ones,
so a retried encode never builds on a partial prior run and the unlink
frees the inodes instead of leaving open fds serving old bytes.
This commit is contained in:
Chris Lu
2026-06-08 16:36:21 -07:00
parent d89a5904e7
commit d069557951
3 changed files with 86 additions and 0 deletions
+25
View File
@@ -88,6 +88,12 @@ func (vs *VolumeServer) VolumeEcShardsGenerate(ctx context.Context, req *volume_
os.Remove(erasure_coding.BitrotSidecarPath(baseFileName, 0))
}()
// Wipe any EC artifacts from a prior encode so a retry never mixes two runs.
// Evict the in-memory EcVolume first so the unlink frees the inodes instead
// of leaving open fds serving the old bytes. Scans the cap for custom ratios.
vs.store.UnloadEcVolume(needle.VolumeId(req.VolumeId))
removeStaleEcArtifacts(baseFileName, v.IndexFileName(), erasure_coding.MaxShardCount)
// IMPORTANT: Generate .ecx BEFORE EC shards to prevent a race condition.
// If .ecx were generated after EC shards, any write (e.g. from WriteNeedleBlob
// during replica sync) between the two steps would add entries to .idx that
@@ -482,6 +488,25 @@ func deleteEcShardIdsForEachLocation(bName string, location *storage.DiskLocatio
return nil
}
// removeStaleEcArtifacts deletes the shard, index, journal, and bitrot sidecar
// files of a prior encode so a fresh encode never mixes runs. total is the
// shard-id range to scan (pass the cap for custom ratios). Does not touch the
// source .dat/.idx/.vif.
func removeStaleEcArtifacts(dataBaseFileName, indexBaseFileName string, total int) {
for i := 0; i < total; i++ {
os.Remove(dataBaseFileName + erasure_coding.ToExt(i))
}
// .ecx/.ecj/.ecsum may sit in either dir depending on -dir.idx; clear both.
os.Remove(indexBaseFileName + ".ecx")
os.Remove(indexBaseFileName + ".ecj")
removeBitrotSidecars(indexBaseFileName)
if dataBaseFileName != indexBaseFileName {
os.Remove(dataBaseFileName + ".ecx")
os.Remove(dataBaseFileName + ".ecj")
removeBitrotSidecars(dataBaseFileName)
}
}
// removeBitrotSidecars removes the legacy <base>.ecsum and any versioned
// <base>.ecsum.v<N> sidecars. Best-effort; logs nothing on absence.
func removeBitrotSidecars(baseFilename string) {
@@ -64,6 +64,59 @@ func TestCheckEcVolumeStatusCountOnlyDataShards(t *testing.T) {
}
}
// TestRemoveStaleEcArtifacts: a fresh encode deletes every prior EC artifact
// (incl. shard ids beyond the default ratio and versioned bitrot sidecars)
// while leaving the source .dat/.idx/.vif untouched.
func TestRemoveStaleEcArtifacts(t *testing.T) {
tempDir := t.TempDir()
dataDir := filepath.Join(tempDir, "data")
idxDir := filepath.Join(tempDir, "idx")
for _, d := range []string{dataDir, idxDir} {
if err := os.MkdirAll(d, 0o755); err != nil {
t.Fatalf("mkdir %s: %v", d, err)
}
}
const baseName = "7"
dataBase := filepath.Join(dataDir, baseName)
idxBase := filepath.Join(idxDir, baseName)
// EC artifacts that must be removed, including a shard id past the default
// 10+4 ratio (proves the MaxShardCount scan) and a versioned sidecar.
var ecFiles []string
for _, id := range []int{0, 9, 13, 20, erasure_coding.MaxShardCount - 1} {
ecFiles = append(ecFiles, dataBase+erasure_coding.ToExt(id))
}
ecFiles = append(ecFiles,
dataBase+".ecx", dataBase+".ecj",
idxBase+".ecx", idxBase+".ecj",
dataBase+erasure_coding.BitrotSidecarExt, // .ecsum (generation 0)
dataBase+erasure_coding.BitrotSidecarExt+".v2", // versioned sidecar
)
// Source files that must survive — the authoritative input for the encode.
srcFiles := []string{dataBase + ".dat", idxBase + ".idx", dataBase + ".vif"}
for _, f := range append(append([]string{}, ecFiles...), srcFiles...) {
if err := os.WriteFile(f, []byte("x"), 0o644); err != nil {
t.Fatalf("create %s: %v", f, err)
}
}
removeStaleEcArtifacts(dataBase, idxBase, erasure_coding.MaxShardCount)
for _, f := range ecFiles {
if util.FileExists(f) {
t.Errorf("expected EC artifact removed: %s", f)
}
}
for _, f := range srcFiles {
if !util.FileExists(f) {
t.Errorf("expected source file preserved: %s", f)
}
}
}
// TestVolumeEcShardsInfo_AggregatesAcrossDisks pins the multi-disk path:
// when a volume server mounts EC shards for the same volume on more than
// one disk (each disk holds its own EcVolume entry — Store.FindEcVolume
+8
View File
@@ -349,6 +349,14 @@ func (s *Store) DestroyEcVolume(vid needle.VolumeId) {
}
}
// UnloadEcVolume drops any in-memory EcVolume for vid from every disk and closes
// its fds without deleting files, so a following unlink frees the inodes.
func (s *Store) UnloadEcVolume(vid needle.VolumeId) {
for _, location := range s.Locations {
location.unloadEcVolume(vid)
}
}
func (s *Store) ReadEcShardNeedle(vid needle.VolumeId, n *needle.Needle, onReadSizeFn func(size types.Size)) (int, error) {
for _, location := range s.Locations {
if localEcVolume, found := location.FindEcVolume(vid); found {