From 675020b34223ba36cacfad0731e4a011dd7b6aab Mon Sep 17 00:00:00 2001 From: Jaehoon Kim Date: Wed, 27 May 2026 12:47:53 +0900 Subject: [PATCH] fix(filer.sync): validate chunk size in FilerSink to prevent 0-byte propagation (#9701) * fix(filer.sync): validate chunk size in FilerSink to prevent 0-byte propagation FilerSink.fetchAndWrite previously trusted the source response and the upload result blindly: a 200 OK / Content-Length: 0 reply from a broken source volume was happily uploaded as a 0-byte needle to the destination, and the destination filer metadata was then written with the source chunk size. The result was permanent silent corruption -- ls shows the file at its original size but reads fail with EIO. Add two cheap defenses inside fetchAndWrite: 1. After assembling fullData, compare its length against sourceChunk.Size. 2. After a successful upload, compare uploadResult.Size against sourceChunk.Size. Both checks wrap a new sentinel errChunkSizeMismatch that the retry callback recognizes and refuses to retry -- needle.size=0 on disk is a persistent state, not a transient network error, so the sync should stop loudly on the affected entry instead of looping or, worse, silently propagating it. Tests: * TestValidateReplicatedChunkSize -- table-driven coverage of healthy, legitimately empty, zero-byte read, short read, and truncated upload cases. * TestFetchAndWriteRejectsZeroByteSource -- end-to-end: an httptest source that returns 200 OK with an empty body must cause fetchAndWrite to return errChunkSizeMismatch after exactly one source hit (fail fast, no retry storm). Co-Authored-By: Claude Opus 4.7 (1M context) * filer.sync: bubble size-mismatch past CreateEntry/UpdateEntry Three follow-ups on the chunk-size validation: - Use %w in replicateOneChunk so the errChunkSizeMismatch sentinel survives the wrap and reaches errors.Is callers up the stack. - In FilerSink.CreateEntry/UpdateEntry, surface errChunkSizeMismatch instead of warning-and-nil. Other errors (deleted source chunk, transient network) keep the existing swallow so a hiccup doesn't stall the stream. - Drop validateReplicatedUploadSize: uploadResult.Size is set client-side from the same len(fullData) we already validated pre-upload, so the second check can't fail. Test: scope the RetryWaitTime override to the one test that needs it, add a regression that locks in the errors.Is chain through replicateChunks. --------- Co-authored-by: Claude Opus 4.7 (1M context) Co-authored-by: Chris Lu --- .../replication/sink/filersink/fetch_write.go | 26 ++- .../sink/filersink/fetch_write_test.go | 203 ++++++++++++++++++ weed/replication/sink/filersink/filer_sink.go | 12 +- 3 files changed, 239 insertions(+), 2 deletions(-) diff --git a/weed/replication/sink/filersink/fetch_write.go b/weed/replication/sink/filersink/fetch_write.go index dd1c3ffb5..b7b7176b8 100644 --- a/weed/replication/sink/filersink/fetch_write.go +++ b/weed/replication/sink/filersink/fetch_write.go @@ -120,7 +120,7 @@ func (fs *FilerSink) replicateOneChunk(sourceChunk *filer_pb.FileChunk, path str fileId, err := fs.fetchAndWrite(sourceChunk, path, sourceMtime) if err != nil { - return nil, fmt.Errorf("copy %s: %v", sourceChunk.GetFileIdString(), err) + return nil, fmt.Errorf("copy %s: %w", sourceChunk.GetFileIdString(), err) } return &filer_pb.FileChunk{ @@ -292,6 +292,10 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string, fullData = data } + if err := validateReplicatedReadSize(sourceChunk, len(fullData)); err != nil { + return err + } + transferStatus.mu.Lock() transferStatus.BytesReceived = int64(len(fullData)) transferStatus.Status = "uploading" @@ -335,6 +339,14 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string, fileId = currentFileId return nil }, func(retryErr error) (shouldContinue bool) { + if errors.Is(retryErr, errChunkSizeMismatch) { + glog.V(0).Infof("permanent size mismatch replicating %s for %s: %v", + sourceChunk.GetFileIdString(), path, retryErr) + transferStatus.mu.Lock() + transferStatus.LastErr = retryErr.Error() + transferStatus.mu.Unlock() + return false + } if fs.hasSourceNewerVersion(path, sourceMtime) { glog.V(1).Infof("skip retrying stale source %s for %s: %v", sourceChunk.GetFileIdString(), path, retryErr) return false @@ -388,6 +400,18 @@ func isEofError(err error) bool { return errors.Is(err, io.ErrUnexpectedEOF) || errors.Is(err, io.EOF) } +// errChunkSizeMismatch is a permanent (non-retriable) replication failure. +var errChunkSizeMismatch = errors.New("chunk size mismatch") + +func validateReplicatedReadSize(sourceChunk *filer_pb.FileChunk, readSize int) error { + if uint64(readSize) != sourceChunk.Size { + return fmt.Errorf("%w: read %s got %d bytes, source metadata says %d", + errChunkSizeMismatch, sourceChunk.GetFileIdString(), + readSize, sourceChunk.Size) + } + return nil +} + func (fs *FilerSink) buildUploadUrl(host, fileId string) string { if fs.writeChunkByFiler { return fmt.Sprintf("http://%s/?proxyChunkId=%s", fs.address, fileId) diff --git a/weed/replication/sink/filersink/fetch_write_test.go b/weed/replication/sink/filersink/fetch_write_test.go index 3238d8fae..90e5fee72 100644 --- a/weed/replication/sink/filersink/fetch_write_test.go +++ b/weed/replication/sink/filersink/fetch_write_test.go @@ -1,12 +1,27 @@ package filersink import ( + "errors" + "net/http" + "net/http/httptest" + "os" + "strings" + "sync/atomic" "testing" + "time" + "github.com/seaweedfs/seaweedfs/weed/operation" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/replication/source" "github.com/seaweedfs/seaweedfs/weed/util" + util_http "github.com/seaweedfs/seaweedfs/weed/util/http" ) +func TestMain(m *testing.M) { + util_http.InitGlobalHttpClient() + os.Exit(m.Run()) +} + func TestTargetPathToSourcePath(t *testing.T) { tests := []struct { name string @@ -77,3 +92,191 @@ func TestTargetPathToSourcePath(t *testing.T) { }) } } + +// FilerSink must reject chunks whose received byte count disagrees with the +// source filer metadata, instead of silently writing 0-byte needles with the +// source size in the destination metadata. +func TestValidateReplicatedChunkSize(t *testing.T) { + const fid = "74,047d16a94aa581" + + tests := []struct { + name string + expectedSize uint64 + readSize int + wantErr bool + }{ + { + name: "healthy", + expectedSize: 5171, + readSize: 5171, + wantErr: false, + }, + { + name: "legitimately empty file", + expectedSize: 0, + readSize: 0, + wantErr: false, + }, + { + name: "zero-byte read for non-empty source", + expectedSize: 5171, + readSize: 0, + wantErr: true, + }, + { + name: "short read", + expectedSize: 5171, + readSize: 100, + wantErr: true, + }, + { + name: "over-read (server returned more than metadata)", + expectedSize: 5171, + readSize: 8192, + wantErr: true, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + chunk := &filer_pb.FileChunk{FileId: fid, Size: tc.expectedSize} + + gotErr := validateReplicatedReadSize(chunk, tc.readSize) + + if tc.wantErr { + if gotErr == nil { + t.Fatalf("expected error, got nil (read=%d expected=%d)", + tc.readSize, tc.expectedSize) + } + if !errors.Is(gotErr, errChunkSizeMismatch) { + t.Fatalf("expected errChunkSizeMismatch, got %v", gotErr) + } + if !strings.Contains(gotErr.Error(), fid) { + t.Fatalf("error %q does not mention chunk id %q", gotErr, fid) + } + return + } + if gotErr != nil { + t.Fatalf("unexpected read-size error: %v", gotErr) + } + }) + } +} + +// End-to-end regression : +// a source volume that responds 200 OK with Content-Length: 0 +// for a chunk that filer metadata claims is 5171 bytes must be rejected +// by fetchAndWrite with a (non-retriable) size mismatch error, +// instead of being silently propagated to the destination as a 0-byte needle. +func TestFetchAndWriteRejectsZeroByteSource(t *testing.T) { + const fid = "74,047d16a94aa581" + const expectedSize uint64 = 5171 + + // Shorten retry backoff so a fail-fast test that briefly enters the retry + // loop doesn't pay the production 1s+ wait. Scoped to this test so any + // future test in the package keeps the production constant. + prevRetryWaitTime := util.RetryWaitTime + util.RetryWaitTime = 100 * time.Millisecond + t.Cleanup(func() { util.RetryWaitTime = prevRetryWaitTime }) + + var hits atomic.Int32 + sourceServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + hits.Add(1) + w.Header().Set("Content-Type", "application/octet-stream") + w.WriteHeader(http.StatusOK) + // Intentionally write no body — mimic the buggy volume response. + })) + defer sourceServer.Close() + + serverAddr := strings.TrimPrefix(sourceServer.URL, "http://") + + filerSrc := &source.FilerSource{} + if err := filerSrc.DoInitialize(serverAddr, serverAddr, "/", true); err != nil { + t.Fatalf("filerSource.DoInitialize: %v", err) + } + + fs := &FilerSink{ + filerSource: filerSrc, + address: serverAddr, + dir: "/dst", + executor: util.NewLimitedConcurrentExecutor(1), + } + fs.SetUploader(operation.NewUploaderWithHttpClient(http.DefaultClient)) + + sourceChunk := &filer_pb.FileChunk{ + FileId: fid, + Size: expectedSize, + } + + done := make(chan struct { + fileId string + err error + }, 1) + go func() { + gotFileId, gotErr := fs.fetchAndWrite(sourceChunk, "/dst/index.bin", 0) + done <- struct { + fileId string + err error + }{gotFileId, gotErr} + }() + + select { + case result := <-done: + if result.err == nil { + t.Fatalf("expected size mismatch error, got nil (fileId=%q)", result.fileId) + } + if !errors.Is(result.err, errChunkSizeMismatch) { + t.Fatalf("expected errChunkSizeMismatch, got %v", result.err) + } + if !strings.Contains(result.err.Error(), "5171") { + t.Fatalf("error %q does not mention expected size 5171", result.err) + } + if !strings.Contains(result.err.Error(), fid) { + t.Fatalf("error %q does not mention chunk id %q", result.err, fid) + } + if h := hits.Load(); h != 1 { + t.Fatalf("expected exactly 1 source hit (fail-fast), got %d", h) + } + case <-time.After(5 * time.Second): + t.Fatalf("fetchAndWrite did not return within 5s (retry loop not aborted on size mismatch); hits=%d", hits.Load()) + } +} + +// Lock in that the errChunkSizeMismatch sentinel survives the wrap in +// replicateOneChunk + pass-through in util.Retry, so filer_sink.go's +// errors.Is check actually fires. +func TestReplicateChunksPreservesSizeMismatchSentinel(t *testing.T) { + const fid = "74,047d16a94aa581" + const expectedSize uint64 = 5171 + + sourceServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/octet-stream") + w.WriteHeader(http.StatusOK) + })) + defer sourceServer.Close() + + serverAddr := strings.TrimPrefix(sourceServer.URL, "http://") + + filerSrc := &source.FilerSource{} + if err := filerSrc.DoInitialize(serverAddr, serverAddr, "/", true); err != nil { + t.Fatalf("filerSource.DoInitialize: %v", err) + } + + fs := &FilerSink{ + filerSource: filerSrc, + address: serverAddr, + dir: "/dst", + executor: util.NewLimitedConcurrentExecutor(1), + } + fs.SetUploader(operation.NewUploaderWithHttpClient(http.DefaultClient)) + + sourceChunks := []*filer_pb.FileChunk{{FileId: fid, Size: expectedSize}} + + _, err := fs.replicateChunks(nil, sourceChunks, "/dst/index.bin", 0) + if err == nil { + t.Fatal("expected error from replicateChunks, got nil") + } + if !errors.Is(err, errChunkSizeMismatch) { + t.Fatalf("error chain broken: errors.Is(err, errChunkSizeMismatch) = false; got %v", err) + } +} diff --git a/weed/replication/sink/filersink/filer_sink.go b/weed/replication/sink/filersink/filer_sink.go index c93942fe9..73ad57e1e 100644 --- a/weed/replication/sink/filersink/filer_sink.go +++ b/weed/replication/sink/filersink/filer_sink.go @@ -2,6 +2,7 @@ package filersink import ( "context" + "errors" "fmt" "math" "sync" @@ -191,7 +192,12 @@ func (fs *FilerSink) CreateEntry(key string, entry *filer_pb.Entry, signatures [ replicatedChunks, err := fs.replicateChunks(context.Background(), entry.GetChunks(), key, getEntryMtime(entry)) if err != nil { - // only warning here since the source chunk may have been deleted already + // Don't swallow size-mismatch: source bytes disagree with source + // metadata, so committing would propagate corruption silently. + if errors.Is(err, errChunkSizeMismatch) { + glog.Errorf("refuse to replicate entry with corrupt chunk %s: %v", key, err) + return err + } glog.Warningf("replicate entry chunks %s: %v", key, err) return nil } @@ -274,6 +280,10 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent // replicate the chunks that are new in the source replicatedChunks, err := fs.replicateChunks(context.Background(), newChunks, key, getEntryMtime(newEntry)) if err != nil { + if errors.Is(err, errChunkSizeMismatch) { + glog.Errorf("refuse to replicate entry with corrupt chunk %s: %v", key, err) + return true, err + } glog.Warningf("replicate entry chunks %s: %v", key, err) return true, nil }