From 9d1c24d80d220bf31e90588a2951b52e2b6adbdb Mon Sep 17 00:00:00 2001 From: ssshr-66 Date: Mon, 5 Oct 2026 09:12:20 +0800 Subject: [PATCH] [Filer] Support append to inline small files (#11591) * fix 11586 * Update filer_server_handlers_write_autochunk.go * filer: fix inline append races, empty files, and stale ETags Serialize the append read-modify-write on the entry lock so concurrent appends merge instead of losing content, keep small appends to empty files inline, tolerate legacy entries whose metadata size differs from their content, and set the entry digest so appended inline files keep a real ETag. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Chris Lu Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../filer_server_handlers_write_autochunk.go | 254 +++- ...lers_write_autochunk_inline_append_test.go | 1060 +++++++++++++++++ .../filer_server_handlers_write_upload.go | 51 + 3 files changed, 1340 insertions(+), 25 deletions(-) create mode 100644 weed/server/filer_server_handlers_write_autochunk_inline_append_test.go diff --git a/weed/server/filer_server_handlers_write_autochunk.go b/weed/server/filer_server_handlers_write_autochunk.go index 28b14893e..2c86affe2 100644 --- a/weed/server/filer_server_handlers_write_autochunk.go +++ b/weed/server/filer_server_handlers_write_autochunk.go @@ -13,6 +13,7 @@ import ( "strings" "time" + "github.com/seaweedfs/seaweedfs/weed/cluster" "github.com/seaweedfs/seaweedfs/weed/filer" "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/operation" @@ -20,6 +21,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/storage/needle" "github.com/seaweedfs/seaweedfs/weed/util" "github.com/seaweedfs/seaweedfs/weed/util/constants" + "google.golang.org/protobuf/proto" ) func (fs *FilerServer) autoChunk(ctx context.Context, w http.ResponseWriter, r *http.Request, contentLength int64, so *operation.StorageOption) { @@ -98,7 +100,7 @@ func (fs *FilerServer) doPostAutoChunk(ctx context.Context, w http.ResponseWrite buf := bufPool.Get().(*bytes.Buffer) buf.Reset() buf.ReadFrom(part1) - filerResult, replyerr = fs.saveMetaData(ctx, r, fileName, contentType, so, nil, nil, 0, buf.Bytes()) + filerResult, replyerr, _ = fs.saveMetaData(ctx, r, fileName, contentType, so, chunkSize, nil, nil, 0, buf.Bytes()) bufPool.Put(buf) return } @@ -114,9 +116,10 @@ func (fs *FilerServer) doPostAutoChunk(ctx context.Context, w http.ResponseWrite fs.filer.DeleteUncommittedChunks(ctx, fileChunks) return nil, nil, errors.New(constants.ErrMsgBadDigest) } - filerResult, replyerr = fs.saveMetaData(ctx, r, fileName, contentType, so, md5bytes, fileChunks, chunkOffset, smallContent) + var uncommittedChunks []*filer_pb.FileChunk + filerResult, replyerr, uncommittedChunks = fs.saveMetaData(ctx, r, fileName, contentType, so, chunkSize, md5bytes, fileChunks, chunkOffset, smallContent) if replyerr != nil { - fs.filer.DeleteUncommittedChunks(ctx, fileChunks) + fs.filer.DeleteUncommittedChunks(ctx, uncommittedChunks) } return @@ -146,9 +149,10 @@ func (fs *FilerServer) doPutAutoChunk(ctx context.Context, w http.ResponseWriter fs.filer.DeleteUncommittedChunks(ctx, fileChunks) return nil, nil, errors.New(constants.ErrMsgBadDigest) } - filerResult, replyerr = fs.saveMetaData(ctx, r, fileName, contentType, so, md5bytes, fileChunks, chunkOffset, smallContent) + var uncommittedChunks []*filer_pb.FileChunk + filerResult, replyerr, uncommittedChunks = fs.saveMetaData(ctx, r, fileName, contentType, so, chunkSize, md5bytes, fileChunks, chunkOffset, smallContent) if replyerr != nil { - fs.filer.DeleteUncommittedChunks(ctx, fileChunks) + fs.filer.DeleteUncommittedChunks(ctx, uncommittedChunks) } return @@ -230,7 +234,8 @@ func (fs *FilerServer) fixFilePath(ctx context.Context, r *http.Request, fileNam return fullPath } -func (fs *FilerServer) saveMetaData(ctx context.Context, r *http.Request, fileName string, contentType string, so *operation.StorageOption, md5bytes []byte, fileChunks []*filer_pb.FileChunk, chunkOffset int64, content []byte) (filerResult *FilerPostResult, replyerr error) { +func (fs *FilerServer) saveMetaData(ctx context.Context, r *http.Request, fileName string, contentType string, so *operation.StorageOption, chunkSize int32, md5bytes []byte, fileChunks []*filer_pb.FileChunk, chunkOffset int64, content []byte) (filerResult *FilerPostResult, replyerr error, uncommittedChunks []*filer_pb.FileChunk) { + uncommittedChunks = fileChunks // detect file mode modeStr := r.URL.Query().Get("mode") @@ -252,33 +257,157 @@ func (fs *FilerServer) saveMetaData(ctx context.Context, r *http.Request, fileNa isAppend := isAppend(r) isOffsetWrite := len(fileChunks) > 0 && fileChunks[0].Offset > 0 + var existingEntry *filer.Entry + var existingSnapshot *filer_pb.Entry + var distributedLock *cluster.LiveLock // when it is an append if isAppend || isOffsetWrite { - existingEntry, findErr := fs.filer.FindEntry(ctx, util.FullPath(path)) - if findErr != nil && findErr != filer_pb.ErrNotFound { - glog.V(0).InfofCtx(ctx, "failing to find %s: %v", path, findErr) - } - entry = existingEntry - } - if entry != nil { - entry.Mtime = time.Now() - entry.Md5 = nil - // adjust chunk offsets - if isAppend { - for _, chunk := range fileChunks { - chunk.Offset += int64(entry.FileSize) + if fs.filer.Dlm != nil && len(fs.filer.Dlm.LockRing.GetSnapshot()) > 1 { + lockClient := cluster.NewLockClient(fs.grpcDialOption, fs.option.Host) + distributedLock = lockClient.NewBlockingLongLivedLock(path, string(fs.option.Host), 0) + if distributedLock == nil { + replyerr = fmt.Errorf("failed to acquire lock for %s; retry the request", path) + return } - entry.FileSize += uint64(chunkOffset) + defer distributedLock.Stop() } - newChunks = append(entry.GetChunks(), fileChunks...) + if fs.entryLockTable != nil { + pathLock := fs.entryLockTable.AcquireLock("appendEntry", util.FullPath(path), util.ExclusiveLock) + defer fs.entryLockTable.ReleaseLock(util.FullPath(path), pathLock) + } + var findErr error + existingEntry, findErr = fs.filer.FindEntry(ctx, util.FullPath(path)) + if findErr != nil && !errors.Is(findErr, filer_pb.ErrNotFound) { + glog.V(0).InfofCtx(ctx, "failing to find %s: %v", path, findErr) + replyerr = fmt.Errorf("find entry %q: %w", path, findErr) + return + } + if existingEntry != nil { + existingSnapshot = existingEntry.ToProtoEntry() + entry = cloneEntryForAppend(existingEntry) + } + } - // TODO - if len(entry.Content) > 0 { - replyerr = fmt.Errorf("append to small file is not supported yet") + inlineAppendHandled := false + inlineAppendConverted := false + if entry != nil { + inlineCandidate := len(entry.Content) > 0 || + (isAppend && entry.FileSize == 0 && entry.Remote == nil && len(entry.HardLinkId) == 0 && len(entry.GetChunks()) == 0) + if isAppend && !so.SaveInside && content != nil && !inlineCandidate { + replyerr = fmt.Errorf("inline file changed while preparing append; retry the request") return } + if inlineCandidate { + if isOffsetWrite { + // TODO: support inline offset writes separately from append semantics. + replyerr = fmt.Errorf("offset write to inline small file is not supported yet") + return + } + if !isAppend || so.SaveInside || entry.Remote != nil || len(entry.HardLinkId) != 0 || len(entry.GetChunks()) != 0 { + replyerr = fmt.Errorf("append to inline content with this storage mode is not supported") + return + } + if entry.FileSize > uint64(len(entry.Content)) { + replyerr = fmt.Errorf("inline file %q has inconsistent size: metadata=%d content=%d", path, entry.FileSize, len(entry.Content)) + return + } + + entry.Mtime = time.Now() + entry.Md5 = nil + oldContentSize := int64(len(entry.Content)) + + switch { + case content != nil: + if int64(len(content)) != chunkOffset || len(fileChunks) != 0 { + replyerr = fmt.Errorf("invalid buffered inline append: size=%d offset=%d chunks=%d", len(content), chunkOffset, len(fileChunks)) + return + } + combinedSize := oldContentSize + int64(len(content)) + if len(content) == 0 || (fs.option.SaveToFilerLimit > 0 && combinedSize < fs.option.SaveToFilerLimit) { + combinedContent := make([]byte, 0, combinedSize) + combinedContent = append(combinedContent, entry.Content...) + combinedContent = append(combinedContent, content...) + entry.Content = combinedContent + entry.Chunks = nil + entry.FileSize = uint64(combinedSize) + entry.Md5 = util.Md5(combinedContent) + newChunks = nil + inlineAppendHandled = true + } else { + // The file grew after the upload path chose inline buffering. Promote + // the complete current file and this append as one chunk stream. + combinedContent := make([]byte, 0, combinedSize) + combinedContent = append(combinedContent, entry.Content...) + combinedContent = append(combinedContent, content...) + fileChunks, _, chunkOffset, replyerr, _ = fs.uploadReaderToChunks(ctx, r, bytes.NewReader(combinedContent), 0, chunkSize, fileName, contentType, true, so) + if replyerr != nil { + return + } + if chunkOffset != combinedSize { + uncommittedChunks = fileChunks + replyerr = fmt.Errorf("promoted inline append size mismatch: uploaded=%d expected=%d", chunkOffset, combinedSize) + return + } + entry.Content = nil + entry.FileSize = uint64(chunkOffset) + newChunks = fileChunks + uncommittedChunks = fileChunks + inlineAppendHandled = true + inlineAppendConverted = true + } + case len(fileChunks) > 0: + prefixChunks, _, prefixOffset, uploadErr, _ := fs.uploadReaderToChunks(ctx, r, bytes.NewReader(entry.Content), 0, chunkSize, fileName, contentType, true, so) + if uploadErr != nil { + replyerr = uploadErr + return + } + uncommittedChunks = append(append([]*filer_pb.FileChunk(nil), fileChunks...), prefixChunks...) + if prefixOffset != oldContentSize { + replyerr = fmt.Errorf("converted inline content size mismatch: uploaded=%d expected=%d", prefixOffset, oldContentSize) + return + } + for _, chunk := range fileChunks { + chunk.Offset += oldContentSize + } + newChunks = append(append([]*filer_pb.FileChunk(nil), prefixChunks...), fileChunks...) + entry.Content = nil + entry.FileSize = uint64(oldContentSize) + uint64(chunkOffset) + inlineAppendHandled = true + inlineAppendConverted = true + uncommittedChunks = newChunks + case chunkOffset == 0: + // An empty append must not migrate an inline file, even if the + // current threshold is disabled or lower than the file's size. + entry.FileSize = uint64(oldContentSize) + entry.Md5 = util.Md5(entry.Content) + newChunks = nil + inlineAppendHandled = true + default: + replyerr = fmt.Errorf("inline append has no content or uploaded chunks") + return + } + } + + if !inlineAppendHandled { + entry.Mtime = time.Now() + entry.Md5 = nil + // New request chunks start at offset zero; append them after the + // existing chunk-backed file. + if isAppend { + for _, chunk := range fileChunks { + chunk.Offset += int64(entry.FileSize) + } + entry.FileSize += uint64(chunkOffset) + } + newChunks = append(append([]*filer_pb.FileChunk(nil), entry.GetChunks()...), fileChunks...) + } + } else { + if isAppend && !so.SaveInside && content != nil { + replyerr = fmt.Errorf("inline append target disappeared while uploading; retry the request") + return + } glog.V(4).InfolnCtx(ctx, "saving", path) newChunks = fileChunks entry = &filer.Entry{ @@ -304,13 +433,22 @@ func (fs *FilerServer) saveMetaData(ctx context.Context, r *http.Request, fileNa glog.V(0).InfofCtx(ctx, "merge chunks %s: %v", r.RequestURI, replyerr) mergedChunks = newChunks } + if inlineAppendConverted { + uncommittedChunks = mergedChunks + } // maybe compact entry chunks mergedChunks, replyerr = filer.MaybeManifestize(fs.saveAsChunk(ctx, so), fs.filer.DeleteChunksNotRecursive, mergedChunks) if replyerr != nil { glog.V(0).InfofCtx(ctx, "manifestize %s: %v", r.RequestURI, replyerr) + if inlineAppendConverted { + uncommittedChunks = mergedChunks + } return } + if inlineAppendConverted { + uncommittedChunks = mergedChunks + } entry.Chunks = mergedChunks if isOffsetWrite { entry.Md5 = nil @@ -339,13 +477,79 @@ func (fs *FilerServer) saveMetaData(ctx context.Context, r *http.Request, fileNa } } + if isAppend || isOffsetWrite { + if distributedLock != nil && !distributedLock.IsLocked() { + replyerr = fmt.Errorf("lost distributed lock for %s; retry the request", path) + return + } + freshEntry, freshErr := fs.filer.FindEntry(ctx, util.FullPath(path)) + if freshErr != nil && !errors.Is(freshErr, filer_pb.ErrNotFound) { + replyerr = fmt.Errorf("failed to recheck entry %q before write: %w", path, freshErr) + return + } + sameEntry := freshEntry == nil && existingSnapshot == nil || + freshEntry != nil && existingSnapshot != nil && + proto.Equal(existingSnapshot, freshEntry.ToProtoEntry()) + if !sameEntry { + replyerr = fmt.Errorf("entry %s changed during append; retry the request", path) + return + } + } + dbErr := fs.filer.CreateEntry(context.WithoutCancel(ctx), entry, nil, false, false, nil, skipCheckParentDirEntry(r), so.MaxFileNameLength) if dbErr != nil { replyerr = dbErr filerResult.Error = dbErr.Error() glog.V(0).InfofCtx(ctx, "failing to write %s to filer server : %v", path, dbErr) + if inlineAppendConverted { + uncommittedChunks = fs.inlineAppendChunksNotReferenced(ctx, util.FullPath(path), mergedChunks) + } + } else { + uncommittedChunks = nil } - return filerResult, replyerr + return +} + +func cloneEntryForAppend(entry *filer.Entry) *filer.Entry { + cloned := entry.ShallowClone() + cloned.Attr.Md5 = append([]byte(nil), entry.Md5...) + cloned.Attr.GroupNames = append([]string(nil), entry.GroupNames...) + cloned.Chunks = append([]*filer_pb.FileChunk(nil), entry.GetChunks()...) + cloned.Content = append([]byte(nil), entry.Content...) + cloned.WORMEnforcedAtTsNs = entry.WORMEnforcedAtTsNs + if entry.Extended != nil { + cloned.Extended = make(map[string][]byte, len(entry.Extended)) + for key, value := range entry.Extended { + cloned.Extended[key] = append([]byte(nil), value...) + } + } + return cloned +} + +func (fs *FilerServer) inlineAppendChunksNotReferenced(ctx context.Context, path util.FullPath, chunks []*filer_pb.FileChunk) []*filer_pb.FileChunk { + verifyCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer cancel() + + entry, err := fs.filer.FindEntry(verifyCtx, path) + if err != nil && !errors.Is(err, filer_pb.ErrNotFound) { + glog.V(0).InfofCtx(verifyCtx, "cannot verify inline append chunks for %s after metadata write failure: %v; retaining them", path, err) + return nil + } + + referenced := make(map[string]struct{}) + if entry != nil { + for _, chunk := range entry.GetChunks() { + referenced[chunk.GetFileIdString()] = struct{}{} + } + } + + var unreferenced []*filer_pb.FileChunk + for _, chunk := range chunks { + if _, found := referenced[chunk.GetFileIdString()]; !found { + unreferenced = append(unreferenced, chunk) + } + } + return unreferenced } func (fs *FilerServer) saveAsChunk(ctx context.Context, so *operation.StorageOption) filer.SaveDataAsChunkFunctionType { diff --git a/weed/server/filer_server_handlers_write_autochunk_inline_append_test.go b/weed/server/filer_server_handlers_write_autochunk_inline_append_test.go new file mode 100644 index 000000000..a65d85e57 --- /dev/null +++ b/weed/server/filer_server_handlers_write_autochunk_inline_append_test.go @@ -0,0 +1,1060 @@ +package weed_server + +import ( + "bytes" + "context" + "crypto/md5" + "errors" + "fmt" + "io" + "math/rand" + "mime/multipart" + "net/http" + "net/http/httptest" + "sort" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/cluster" + "github.com/seaweedfs/seaweedfs/weed/filer" + "github.com/seaweedfs/seaweedfs/weed/operation" + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" + "github.com/seaweedfs/seaweedfs/weed/util" + "github.com/seaweedfs/seaweedfs/weed/wdclient" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" +) + +func TestFilerInlineAppendStaysInlineBelowLimit(t *testing.T) { + ctx := context.Background() + const ( + filePath = "/append.txt" + limit = int64(12) + ) + oldContent := []byte("old-") + appendContent := []byte("new-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + originalExtended := map[string][]byte{"X-Existing": []byte("preserved")} + entry := &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + Mode: 0640, + Uid: 123, + Gid: 456, + FileSize: uint64(len(oldContent)), + }, + Extended: originalExtended, + Content: append([]byte(nil), oldContent...), + } + if err := store.InsertEntry(ctx, entry); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs := &FilerServer{ + filer: f, + option: &FilerOption{ + SaveToFilerLimit: limit, + MaxMB: 1, + }, + } + so := &operation.StorageOption{} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", bytes.NewReader(appendContent)) + r.ContentLength = -1 // The inline decision must be based on the stream, not Content-Length. + r.Header.Set("Cache-Control", "max-age=60") + + fileChunks, md5Hash, chunkOffset, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, r.Body, 4, "append.txt", "", -1, so, + ) + if err != nil { + t.Fatalf("prepare append: %v", err) + } + if len(fileChunks) != 0 { + t.Fatalf("inline append uploaded %d chunks, want none", len(fileChunks)) + } + if md5Hash == nil { + t.Fatal("append digest is nil") + } + if chunkOffset != int64(len(appendContent)) { + t.Fatalf("append size = %d, want %d", chunkOffset, len(appendContent)) + } + if !bytes.Equal(smallContent, appendContent) { + t.Fatalf("buffered append = %q, want %q", smallContent, appendContent) + } + + result, err, cleanupChunks := fs.saveMetaData(ctx, r, "append.txt", "", so, 4, md5Hash.Sum(nil), fileChunks, chunkOffset, smallContent) + if err != nil { + t.Fatalf("save appended entry: %v", err) + } + if len(cleanupChunks) != 0 { + t.Fatalf("successful append returned uncommitted chunks: %#v", cleanupChunks) + } + if result == nil || result.Size != int64(len(oldContent)+len(appendContent)) { + t.Fatalf("result = %#v, want size %d", result, len(oldContent)+len(appendContent)) + } + + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + wantContent := append(append([]byte(nil), oldContent...), appendContent...) + if !bytes.Equal(updated.Content, wantContent) { + t.Fatalf("stored content = %q, want %q", updated.Content, wantContent) + } + if len(updated.Chunks) != 0 { + t.Fatalf("stored %d chunks for below-limit append, want inline content", len(updated.Chunks)) + } + if updated.FileSize != uint64(len(wantContent)) { + t.Fatalf("file size = %d, want %d", updated.FileSize, len(wantContent)) + } + if string(updated.Extended["X-Existing"]) != "preserved" { + t.Fatalf("existing extended attribute was lost: %#v", updated.Extended) + } + if string(updated.Extended["Cache-Control"]) != "max-age=60" { + t.Fatalf("request extended attribute was not saved: %#v", updated.Extended) + } + if _, mutated := originalExtended["Cache-Control"]; mutated { + t.Fatalf("the original entry's extended attributes were mutated: %#v", originalExtended) + } + if updated.Mtime.Equal(time.Unix(1, 0)) { + t.Fatal("mtime was not updated") + } +} + +func TestFilerInlineAppendEmptyKeepsInlineContent(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + // The current threshold is below the existing file size. An empty append + // must not migrate the file merely because it is now over that threshold. + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 3, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", bytes.NewReader(nil)) + fileChunks, md5Hash, chunkOffset, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, r.Body, 4, "append.txt", "", 0, &operation.StorageOption{}, + ) + if err != nil { + t.Fatalf("prepare empty append: %v", err) + } + if len(fileChunks) != 0 || smallContent == nil || len(smallContent) != 0 { + t.Fatalf("empty append upload state: chunks=%d content=%#v", len(fileChunks), smallContent) + } + if md5Hash == nil || chunkOffset != 0 { + t.Fatalf("empty append digest/size = %v/%d, want non-nil digest and size 0", md5Hash, chunkOffset) + } + + if _, err, cleanupChunks := fs.saveMetaData(ctx, r, "append.txt", "", &operation.StorageOption{}, 4, md5Hash.Sum(nil), fileChunks, chunkOffset, smallContent); err != nil { + t.Fatalf("save empty append: %v", err) + } else if len(cleanupChunks) != 0 { + t.Fatalf("empty append returned cleanup chunks: %#v", cleanupChunks) + } + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + if !bytes.Equal(updated.Content, oldContent) || len(updated.Chunks) != 0 || updated.FileSize != uint64(len(oldContent)) { + t.Fatalf("empty append changed storage representation: content=%q chunks=%d size=%d", updated.Content, len(updated.Chunks), updated.FileSize) + } +} + +func TestFilerInlineAppendRejectsPromotionWithoutChunkSize(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 8, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", strings.NewReader("new-")) + chunks, _, _, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, r.Body, 0, "append.txt", "", -1, &operation.StorageOption{}, + ) + if err == nil || !strings.Contains(err.Error(), "invalid chunk size 0") { + t.Fatalf("promotion error = %v, want invalid chunk size", err) + } + if len(chunks) != 0 || smallContent != nil { + t.Fatalf("invalid promotion returned data: chunks=%d inline=%q", len(chunks), smallContent) + } + unchanged, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find original entry: %v", err) + } + if !bytes.Equal(unchanged.Content, oldContent) { + t.Fatalf("invalid promotion changed original content: %q", unchanged.Content) + } +} + +func TestFilerInlineAppendMultipartPostStaysInline(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + appendContent := []byte("new-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + + var body bytes.Buffer + multipartWriter := multipart.NewWriter(&body) + part, err := multipartWriter.CreateFormFile("file", "append.txt") + if err != nil { + t.Fatalf("create multipart file part: %v", err) + } + if _, err = part.Write(appendContent); err != nil { + t.Fatalf("write multipart file part: %v", err) + } + if err = multipartWriter.Close(); err != nil { + t.Fatalf("close multipart body: %v", err) + } + r := httptest.NewRequest(http.MethodPost, "/?op=append", &body) + r.Header.Set("Content-Type", multipartWriter.FormDataContentType()) + recorder := httptest.NewRecorder() + result, requestMD5, err := fs.doPostAutoChunk(ctx, recorder, r, 4, int64(body.Len()), &operation.StorageOption{}) + if err != nil { + t.Fatalf("POST inline append: %v", err) + } + if result == nil || result.Size != int64(len(oldContent)+len(appendContent)) { + t.Fatalf("result = %#v, want size %d", result, len(oldContent)+len(appendContent)) + } + wantMD5 := md5.Sum(appendContent) + if !bytes.Equal(requestMD5, wantMD5[:]) { + t.Fatalf("request MD5 = %x, want %x", requestMD5, wantMD5) + } + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + wantContent := append(append([]byte(nil), oldContent...), appendContent...) + if !bytes.Equal(updated.Content, wantContent) || len(updated.Chunks) != 0 { + t.Fatalf("multipart append result: content=%q chunks=%d, want inline %q", updated.Content, len(updated.Chunks), wantContent) + } +} + +func TestFilerInlineOffsetWriteRemainsUnsupported(t *testing.T) { + ctx := context.Background() + const filePath = "/offset.txt" + oldContent := []byte("keep") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?offset=1", bytes.NewReader([]byte("X"))) + _, err, _ := fs.saveMetaData(ctx, r, "offset.txt", "", &operation.StorageOption{}, 4, nil, []*filer_pb.FileChunk{{Offset: 1, Size: 1}}, 2, nil) + if err == nil || !strings.Contains(err.Error(), "offset write to inline small file is not supported") { + t.Fatalf("offset write error = %v, want explicit unsupported error", err) + } + unchanged, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find original entry: %v", err) + } + if !bytes.Equal(unchanged.Content, oldContent) || unchanged.FileSize != uint64(len(oldContent)) { + t.Fatalf("unsupported offset write changed entry: content=%q size=%d", unchanged.Content, unchanged.FileSize) + } +} + +func TestFilerInlineAppendReadErrorKeepsOriginalEntry(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", nil) + fileChunks, _, _, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, &inlineAppendReadError{}, 4, "append.txt", "", -1, &operation.StorageOption{}, + ) + if err == nil || !strings.Contains(err.Error(), "read input: injected read failure") { + t.Fatalf("read error = %v, want propagated input read failure", err) + } + if len(fileChunks) != 0 || smallContent != nil { + t.Fatalf("failed read produced chunks/content: chunks=%d content=%q", len(fileChunks), smallContent) + } + unchanged, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find original entry: %v", err) + } + if !bytes.Equal(unchanged.Content, oldContent) { + t.Fatalf("failed read changed original content: %q", unchanged.Content) + } +} + +type inlineAppendReadError struct { + readData bool +} + +func (r *inlineAppendReadError) Read(p []byte) (int, error) { + if r.readData { + return 0, errors.New("injected read failure") + } + r.readData = true + return copy(p, []byte("new-")), nil +} + +func TestFilerInlineAppendBadDigestKeepsOriginalEntry(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", strings.NewReader("new-")) + r.Header.Set("Content-MD5", "not-the-correct-digest") + _, _, err := fs.doPutAutoChunk(ctx, httptest.NewRecorder(), r, 4, 4, &operation.StorageOption{}) + if err == nil || !strings.Contains(err.Error(), "Content-Md5") { + t.Fatalf("PUT error = %v, want bad Content-MD5 error", err) + } + unchanged, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find original entry: %v", err) + } + if !bytes.Equal(unchanged.Content, oldContent) || len(unchanged.Chunks) != 0 { + t.Fatalf("bad digest changed original entry: content=%q chunks=%d", unchanged.Content, len(unchanged.Chunks)) + } +} + +func TestFilerChunkAppendStillOffsetsNewChunks(t *testing.T) { + ctx := context.Background() + const filePath = "/chunked.txt" + oldChunk := &filer_pb.FileChunk{FileId: "1,00000001", Offset: 0, Size: 4} + newChunk := &filer_pb.FileChunk{FileId: "1,00000002", Offset: 0, Size: 3} + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: 4, + }, + Chunks: []*filer_pb.FileChunk{oldChunk}, + }); err != nil { + t.Fatalf("insert chunk-backed entry: %v", err) + } + + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", strings.NewReader("new")) + result, err, cleanupChunks := fs.saveMetaData(ctx, r, "chunked.txt", "", &operation.StorageOption{}, 4, nil, []*filer_pb.FileChunk{newChunk}, 3, nil) + if err != nil { + t.Fatalf("append to chunk-backed entry: %v", err) + } + if len(cleanupChunks) != 0 { + t.Fatalf("successful chunk append returned cleanup chunks: %#v", cleanupChunks) + } + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + if result == nil || result.Size != 7 || updated.FileSize != 7 || len(updated.Chunks) != 2 { + t.Fatalf("append result=%#v entry-size=%d chunks=%d, want size 7 and two chunks", result, updated.FileSize, len(updated.Chunks)) + } + if updated.Chunks[0].Offset != 0 || updated.Chunks[0].GetFileIdString() != oldChunk.GetFileIdString() || updated.Chunks[1].Offset != 4 || updated.Chunks[1].GetFileIdString() != newChunk.GetFileIdString() { + t.Fatalf("chunk append metadata = %#v, want offsets 0 and 4 with old and new file IDs", updated.Chunks) + } +} + +func TestFilerInlineAppendPromotesAtLimit(t *testing.T) { + for _, tc := range []struct { + name string + appendContent string + failStore bool + failVerify bool + }{ + {name: "exactly reaches limit", appendContent: "new-"}, + {name: "crosses limit", appendContent: "new!!"}, + {name: "metadata write fails", appendContent: "new-", failStore: true}, + {name: "metadata result cannot be verified", appendContent: "new-", failStore: true, failVerify: true}, + } { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + const ( + filePath = "/append.txt" + limit = int64(8) + ) + oldContent := []byte("old-") + appendContent := []byte(tc.appendContent) + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + Mode: 0640, + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs, uploadedChunks := newInlineAppendUploadServer(t, f, limit, 4) + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", bytes.NewReader(appendContent)) + r.ContentLength = -1 + fileChunks, md5Hash, chunkOffset, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, r.Body, 4, "append.txt", "", -1, &operation.StorageOption{}, + ) + if err != nil { + t.Fatalf("prepare append: %v", err) + } + if len(fileChunks) == 0 { + t.Fatal("append reaching the inline limit was not uploaded as chunks") + } + if smallContent != nil { + t.Fatalf("append reaching the inline limit was buffered as inline content: %q", smallContent) + } + if chunkOffset != int64(len(appendContent)) { + t.Fatalf("append size = %d, want %d", chunkOffset, len(appendContent)) + } + wantAppendMD5 := md5.Sum(appendContent) + if got := md5Hash.Sum(nil); !bytes.Equal(got, wantAppendMD5[:]) { + t.Fatalf("append MD5 = %x, want %x", got, wantAppendMD5) + } + + if tc.failVerify { + f.Store = filer.NewFilerStoreWrapper(&inlineAppendUncertainUpdateStore{ + renameTestStore: store, + err: errors.New("injected ambiguous metadata update failure"), + }) + } else if tc.failStore { + f.Store = filer.NewFilerStoreWrapper(&inlineAppendFailUpdateStore{ + renameTestStore: store, + err: errors.New("injected metadata update failure"), + }) + } + + result, err, cleanupChunks := fs.saveMetaData(ctx, r, "append.txt", "", &operation.StorageOption{}, 4, md5Hash.Sum(nil), fileChunks, chunkOffset, smallContent) + if tc.failStore { + if err == nil { + t.Fatal("save appended entry succeeded despite injected metadata failure") + } + if tc.failVerify { + if len(cleanupChunks) != 0 { + t.Fatalf("unverified metadata state returned %d chunks for deletion, want retain all", len(cleanupChunks)) + } + } else if len(cleanupChunks) != len(fileChunks)+1 { + t.Fatalf("cleanup chunk count = %d, want uploaded append chunks %d plus one converted prefix chunk", len(cleanupChunks), len(fileChunks)) + } + for _, chunk := range cleanupChunks { + if _, found := uploadedChunks.get(chunk.GetFileIdString()); !found { + t.Fatalf("cleanup includes chunk %s without uploaded data", chunk.GetFileIdString()) + } + } + unchanged, findErr := store.FindEntry(ctx, util.FullPath(filePath)) + if findErr != nil { + t.Fatalf("find original entry after failed update: %v", findErr) + } + if !bytes.Equal(unchanged.Content, oldContent) || len(unchanged.Chunks) != 0 { + t.Fatalf("failed update changed original entry: content=%q chunks=%d", unchanged.Content, len(unchanged.Chunks)) + } + return + } + if err != nil { + t.Fatalf("save appended entry: %v", err) + } + if len(cleanupChunks) != 0 { + t.Fatalf("successful append returned uncommitted chunks: %#v", cleanupChunks) + } + + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + if len(updated.Content) != 0 { + t.Fatalf("promoted entry still has inline content %q", updated.Content) + } + if len(updated.Chunks) == 0 { + t.Fatal("promoted entry has no chunks") + } + wantContent := append(append([]byte(nil), oldContent...), appendContent...) + if result == nil || result.Size != int64(len(wantContent)) { + t.Fatalf("result = %#v, want size %d", result, len(wantContent)) + } + if updated.FileSize != uint64(len(wantContent)) { + t.Fatalf("file size = %d, want %d", updated.FileSize, len(wantContent)) + } + + sort.Slice(updated.Chunks, func(i, j int) bool { + return updated.Chunks[i].Offset < updated.Chunks[j].Offset + }) + var gotContent []byte + for _, chunk := range updated.Chunks { + if chunk.Offset != int64(len(gotContent)) { + t.Fatalf("chunk %s starts at %d, want contiguous offset %d", chunk.FileId, chunk.Offset, len(gotContent)) + } + data, found := uploadedChunks.get(chunk.GetFileIdString()) + if !found { + t.Fatalf("no uploaded data for chunk %s", chunk.GetFileIdString()) + } + if uint64(len(data)) != chunk.Size { + t.Fatalf("chunk %s metadata size = %d, uploaded size = %d", chunk.FileId, chunk.Size, len(data)) + } + gotContent = append(gotContent, data...) + } + if !bytes.Equal(gotContent, wantContent) { + t.Fatalf("chunk data = %q, want %q", gotContent, wantContent) + } + }) + } +} + +func TestFilerInlineAppendRealThresholdBoundaries(t *testing.T) { + const ( + filePath = "/append-boundary.bin" + limit = 65536 + oldSize = 4096 + ) + for _, finalSize := range []int{limit - 1, limit, limit + 1} { + t.Run(fmt.Sprintf("final_size_%d", finalSize), func(t *testing.T) { + ctx := context.Background() + oldContent := inlineAppendTestBytes(oldSize, 1) + appendContent := inlineAppendTestBytes(finalSize-oldSize, 2) + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: oldContent, + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs, uploaded := newInlineAppendUploadServer(t, f, limit, 1) + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", bytes.NewReader(appendContent)) + r.ContentLength = -1 + fileChunks, md5Hash, chunkOffset, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, r.Body, limit, "append-boundary.bin", "", -1, &operation.StorageOption{}, + ) + if err != nil { + t.Fatalf("prepare append: %v", err) + } + if md5Hash == nil || chunkOffset != int64(len(appendContent)) { + t.Fatalf("append digest/size = %v/%d, want non-nil digest and %d", md5Hash, chunkOffset, len(appendContent)) + } + + result, saveErr, cleanupChunks := fs.saveMetaData(ctx, r, "append-boundary.bin", "", &operation.StorageOption{}, limit, md5Hash.Sum(nil), fileChunks, chunkOffset, smallContent) + if saveErr != nil { + t.Fatalf("save appended entry: %v", saveErr) + } + if len(cleanupChunks) != 0 { + t.Fatalf("successful append returned cleanup chunks: %#v", cleanupChunks) + } + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + if result == nil || result.Size != int64(finalSize) || updated.FileSize != uint64(finalSize) { + t.Fatalf("result=%#v entry size=%d, want %d", result, updated.FileSize, finalSize) + } + + wantContent := append(append([]byte(nil), oldContent...), appendContent...) + if finalSize < limit { + if len(fileChunks) != 0 || smallContent == nil || uploaded.count() != 0 { + t.Fatalf("below-limit path uploaded data: chunks=%d inline=%d volume uploads=%d", len(fileChunks), len(smallContent), uploaded.count()) + } + if !bytes.Equal(updated.Content, wantContent) || len(updated.Chunks) != 0 { + t.Fatalf("below-limit entry: content size=%d chunks=%d", len(updated.Content), len(updated.Chunks)) + } + return + } + + if smallContent != nil || len(updated.Content) != 0 || len(updated.Chunks) == 0 { + t.Fatalf("at/above-limit entry: inline=%d chunks=%d", len(updated.Content), len(updated.Chunks)) + } + sort.Slice(updated.Chunks, func(i, j int) bool { + return updated.Chunks[i].Offset < updated.Chunks[j].Offset + }) + var reconstructed []byte + for _, chunk := range updated.Chunks { + if chunk.Offset != int64(len(reconstructed)) { + t.Fatalf("chunk offset = %d, want %d", chunk.Offset, len(reconstructed)) + } + data, found := uploaded.get(chunk.GetFileIdString()) + if !found { + t.Fatalf("missing Volume data for %s", chunk.GetFileIdString()) + } + if chunk.IsCompressed { + data, err = util.DecompressData(data) + if err != nil { + t.Fatalf("decompress uploaded chunk %s: %v", chunk.GetFileIdString(), err) + } + } + if uint64(len(data)) != chunk.Size { + t.Fatalf("chunk size = %d, want %d", len(data), chunk.Size) + } + reconstructed = append(reconstructed, data...) + } + if !bytes.Equal(reconstructed, wantContent) { + t.Fatalf("reconstructed %d bytes, want %d", len(reconstructed), len(wantContent)) + } + }) + } +} + +func TestFilerInlineAppendEmptyFileStaysInline(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + appendContent := []byte("new-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + }, + }); err != nil { + t.Fatalf("insert empty entry: %v", err) + } + + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", bytes.NewReader(appendContent)) + r.ContentLength = -1 + fileChunks, md5Hash, chunkOffset, err, smallContent := fs.uploadRequestToChunks( + ctx, httptest.NewRecorder(), r, r.Body, 4, "append.txt", "", -1, &operation.StorageOption{}, + ) + if err != nil { + t.Fatalf("prepare append to empty file: %v", err) + } + if len(fileChunks) != 0 || !bytes.Equal(smallContent, appendContent) { + t.Fatalf("append to empty file uploaded instead of inline: chunks=%d inline=%q", len(fileChunks), smallContent) + } + + result, err, cleanupChunks := fs.saveMetaData(ctx, r, "append.txt", "", &operation.StorageOption{}, 4, md5Hash.Sum(nil), fileChunks, chunkOffset, smallContent) + if err != nil { + t.Fatalf("save append to empty file: %v", err) + } + if len(cleanupChunks) != 0 || result == nil || result.Size != int64(len(appendContent)) { + t.Fatalf("append result=%#v cleanup=%d", result, len(cleanupChunks)) + } + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + if !bytes.Equal(updated.Content, appendContent) || len(updated.Chunks) != 0 || updated.FileSize != uint64(len(appendContent)) { + t.Fatalf("append to empty file: content=%q chunks=%d size=%d", updated.Content, len(updated.Chunks), updated.FileSize) + } +} + +func TestFilerInlineAppendKeepsETag(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + appendContent := []byte("new-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", bytes.NewReader(appendContent)) + r.ContentLength = -1 + if _, err, cleanupChunks := fs.saveMetaData(ctx, r, "append.txt", "", &operation.StorageOption{}, 4, nil, nil, int64(len(appendContent)), appendContent); err != nil { + t.Fatalf("save inline append: %v", err) + } else if len(cleanupChunks) != 0 { + t.Fatalf("inline append returned cleanup chunks: %#v", cleanupChunks) + } + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + wantContent := append(append([]byte(nil), oldContent...), appendContent...) + wantMD5 := md5.Sum(wantContent) + if !bytes.Equal(updated.Md5, wantMD5[:]) { + t.Fatalf("entry MD5 = %x, want %x", updated.Md5, wantMD5) + } + if got, want := filer.ETagEntry(updated), fmt.Sprintf("%x", wantMD5); got != want { + t.Fatalf("entry ETag = %q, want %q", got, want) + } +} + +func TestFilerInlineAppendConcurrentKeepsAllData(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + oldContent := []byte("old-") + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + f.Store = filer.NewFilerStoreWrapper(&inlineAppendSlowFindStore{renameTestStore: store, delay: 2 * time.Millisecond}) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: uint64(len(oldContent)), + }, + Content: append([]byte(nil), oldContent...), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs := &FilerServer{ + filer: f, + option: &FilerOption{SaveToFilerLimit: 64, MaxMB: 1}, + entryLockTable: util.NewLockTable[util.FullPath](), + } + markers := []string{"A1", "B2", "C3", "D4", "E5", "F6", "G7", "H8"} + var wg sync.WaitGroup + errs := make(chan error, len(markers)) + for _, marker := range markers { + wg.Add(1) + go func() { + defer wg.Done() + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", strings.NewReader(marker)) + _, err, _ := fs.saveMetaData(ctx, r, "append.txt", "", &operation.StorageOption{}, 4, nil, nil, int64(len(marker)), []byte(marker)) + errs <- err + }() + } + wg.Wait() + close(errs) + for err := range errs { + if err != nil { + t.Fatalf("concurrent append failed: %v", err) + } + } + + updated, err := f.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find updated entry: %v", err) + } + if len(updated.Chunks) != 0 { + t.Fatalf("concurrent appends promoted the entry to %d chunks", len(updated.Chunks)) + } + got := string(updated.Content) + if !strings.HasPrefix(got, string(oldContent)) { + t.Fatalf("concurrent append content = %q, want prefix %q", got, oldContent) + } + for _, marker := range markers { + if !strings.Contains(got, marker) { + t.Fatalf("concurrent append lost marker %q: content = %q", marker, got) + } + } + if updated.FileSize != uint64(len(oldContent)+len(markers)*2) { + t.Fatalf("final size = %d, want %d", updated.FileSize, len(oldContent)+len(markers)*2) + } +} + +type inlineAppendSlowFindStore struct { + *renameTestStore + delay time.Duration +} + +func (s *inlineAppendSlowFindStore) FindEntry(ctx context.Context, path util.FullPath) (*filer.Entry, error) { + time.Sleep(s.delay) + return s.renameTestStore.FindEntry(ctx, path) +} + +type inlineAppendFailUpdateStore struct { + *renameTestStore + err error +} + +func (s *inlineAppendFailUpdateStore) UpdateEntry(context.Context, *filer.Entry) error { + return s.err +} + +type inlineAppendUncertainUpdateStore struct { + *renameTestStore + err error + updateFailed atomic.Bool +} + +func (s *inlineAppendUncertainUpdateStore) UpdateEntry(context.Context, *filer.Entry) error { + s.updateFailed.Store(true) + return s.err +} + +func (s *inlineAppendUncertainUpdateStore) FindEntry(ctx context.Context, path util.FullPath) (*filer.Entry, error) { + if s.updateFailed.Load() { + return nil, errors.New("injected verification lookup failure") + } + return s.renameTestStore.FindEntry(ctx, path) +} + +type inlineAppendVolumeData struct { + mu sync.Mutex + data map[string][]byte +} + +func (d *inlineAppendVolumeData) get(fileID string) ([]byte, bool) { + d.mu.Lock() + defer d.mu.Unlock() + data, ok := d.data[fileID] + return append([]byte(nil), data...), ok +} + +func (d *inlineAppendVolumeData) count() int { + d.mu.Lock() + defer d.mu.Unlock() + return len(d.data) +} + +func inlineAppendTestBytes(length int, seed int64) []byte { + data := make([]byte, length) + _, _ = rand.New(rand.NewSource(seed)).Read(data) + return data +} + +type inlineAppendFakeMaster struct { + master_pb.UnimplementedSeaweedServer + volumeHost string + nextFileID atomic.Uint32 +} + +func (m *inlineAppendFakeMaster) KeepConnected(stream grpc.BidiStreamingServer[master_pb.KeepConnectedRequest, master_pb.KeepConnectedResponse]) error { + if _, err := stream.Recv(); err != nil { + return err + } + if err := stream.Send(&master_pb.KeepConnectedResponse{}); err != nil { + return err + } + <-stream.Context().Done() + return stream.Context().Err() +} + +func (m *inlineAppendFakeMaster) Assign(_ context.Context, _ *master_pb.AssignRequest) (*master_pb.AssignResponse, error) { + n := m.nextFileID.Add(1) + return &master_pb.AssignResponse{ + Fid: fmt.Sprintf("1,%08x", n), + Count: 1, + Location: &master_pb.Location{ + Url: m.volumeHost, + PublicUrl: m.volumeHost, + }, + }, nil +} + +func newInlineAppendUploadServer(t *testing.T, f *filer.Filer, saveToFilerLimit int64, maxMB int) (*FilerServer, *inlineAppendVolumeData) { + t.Helper() + + uploaded := &inlineAppendVolumeData{data: make(map[string][]byte)} + volumeServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + multipartReader, err := r.MultipartReader() + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + part, err := multipartReader.NextPart() + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + data, err := io.ReadAll(part) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + fileID := strings.TrimPrefix(r.URL.Path, "/") + uploaded.mu.Lock() + uploaded.data[fileID] = append([]byte(nil), data...) + uploaded.mu.Unlock() + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusCreated) + _, _ = fmt.Fprintf(w, `{"name":"chunk","size":%d}`, len(data)) + })) + t.Cleanup(volumeServer.Close) + + volumeHost := strings.TrimPrefix(volumeServer.URL, "http://") + masterAddr := startFakeMasterServerForLeaderLookup(t, &inlineAppendFakeMaster{volumeHost: volumeHost}) + dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) + masterClient := wdclient.NewMasterClient( + dialOption, + "test", + cluster.FilerType, + pb.ServerAddress("localhost:0"), + "", + "", + *pb.NewServiceDiscoveryFromMap(map[string]pb.ServerAddress{"master": masterAddr}), + ) + f.MasterClient = masterClient + masterCtx, cancelMaster := context.WithCancel(context.Background()) + t.Cleanup(cancelMaster) + go masterClient.KeepConnectedToMaster(masterCtx) + + waitCtx, cancelWait := context.WithTimeout(context.Background(), 5*time.Second) + defer cancelWait() + if masterClient.GetMaster(waitCtx) == "" { + t.Fatal("fake master did not become available") + } + + return &FilerServer{ + filer: f, + option: &FilerOption{SaveToFilerLimit: saveToFilerLimit, MaxMB: maxMB}, + grpcDialOption: dialOption, + }, uploaded +} + +// inlineAppendDriftingStore commits a competing update right after the read +// that prepares an append, before the recheck ahead of the metadata commit. +type inlineAppendDriftingStore struct { + *renameTestStore + reads int + drifted bool +} + +func (s *inlineAppendDriftingStore) FindEntry(ctx context.Context, path util.FullPath) (*filer.Entry, error) { + entry, err := s.renameTestStore.FindEntry(ctx, path) + if err == nil && entry != nil { + s.reads++ + if s.reads == 2 && !s.drifted { + s.drifted = true + competing := entry.ShallowClone() + competing.Content = append([]byte(nil), append(competing.Content, '!')...) + competing.FileSize = uint64(len(competing.Content)) + if updateErr := s.renameTestStore.UpdateEntry(ctx, competing); updateErr != nil { + return nil, updateErr + } + } + } + return entry, err +} + +func TestFilerInlineAppendRejectsEntryChangedMidWrite(t *testing.T) { + ctx := context.Background() + const filePath = "/append.txt" + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + drifting := &inlineAppendDriftingStore{renameTestStore: store} + f.Store = filer.NewFilerStoreWrapper(drifting) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: 4, + }, + Content: []byte("old-"), + }); err != nil { + t.Fatalf("insert inline entry: %v", err) + } + + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 64, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", strings.NewReader("new")) + if _, err, _ := fs.saveMetaData(ctx, r, "append.txt", "", &operation.StorageOption{}, 4, nil, nil, 3, []byte("new")); err == nil { + t.Fatal("append succeeded despite the entry changing before commit") + } + + unchanged, err := store.FindEntry(ctx, util.FullPath(filePath)) + if err != nil { + t.Fatalf("find entry after rejected append: %v", err) + } + if string(unchanged.Content) != "old-!" { + t.Fatalf("rejected append overwrote the competing write: %q", unchanged.Content) + } +} + +func TestFilerChunkAppendWithExtendedHeaders(t *testing.T) { + ctx := context.Background() + const filePath = "/chunked.txt" + oldChunk := &filer_pb.FileChunk{FileId: "1,00000001", Offset: 0, Size: 4} + newChunk := &filer_pb.FileChunk{FileId: "1,00000002", Offset: 0, Size: 3} + store := newRenameTestStore() + f := newRenameTestFiler(t, store) + if err := store.InsertEntry(ctx, &filer.Entry{ + FullPath: util.FullPath(filePath), + Attr: filer.Attr{ + Mtime: time.Unix(1, 0), + Crtime: time.Unix(1, 0), + FileSize: 4, + }, + Chunks: []*filer_pb.FileChunk{oldChunk}, + Extended: map[string][]byte{"existing": []byte("value")}, + }); err != nil { + t.Fatalf("insert chunk-backed entry: %v", err) + } + + fs := &FilerServer{filer: f, option: &FilerOption{SaveToFilerLimit: 12, MaxMB: 1}} + r := httptest.NewRequest(http.MethodPut, filePath+"?op=append", strings.NewReader("new")) + r.Header.Set("Cache-Control", "max-age=60") + result, err, _ := fs.saveMetaData(ctx, r, "chunked.txt", "", &operation.StorageOption{}, 4, nil, []*filer_pb.FileChunk{newChunk}, 3, nil) + if err != nil { + t.Fatalf("append with extended-attribute header: %v", err) + } + if result == nil || result.Size != 7 { + t.Fatalf("append result=%#v, want size 7", result) + } +} diff --git a/weed/server/filer_server_handlers_write_upload.go b/weed/server/filer_server_handlers_write_upload.go index f2c18043a..f41ba6151 100644 --- a/weed/server/filer_server_handlers_write_upload.go +++ b/weed/server/filer_server_handlers_write_upload.go @@ -5,6 +5,7 @@ import ( "context" "crypto/md5" "encoding/base64" + "errors" "fmt" "hash" "io" @@ -47,6 +48,56 @@ func (fs *FilerServer) uploadRequestToChunks(ctx context.Context, w http.Respons chunkOffset = offsetInt } + if isAppend && !so.SaveInside { + fullPath := fs.fixFilePath(ctx, r, fileName) + entry, findErr := fs.filer.FindEntry(ctx, util.FullPath(fullPath)) + if findErr != nil && !errors.Is(findErr, filer_pb.ErrNotFound) { + return nil, nil, 0, fmt.Errorf("find entry for append %q: %w", fullPath, findErr), nil + } + if findErr == nil && entry != nil && !entry.IsDirectory() && entry.Remote == nil && len(entry.HardLinkId) == 0 && len(entry.GetChunks()) == 0 && (len(entry.Content) > 0 || entry.FileSize == 0) { + if entry.FileSize > uint64(len(entry.Content)) { + return nil, nil, 0, fmt.Errorf("inline file %q has inconsistent size: metadata=%d content=%d", fullPath, entry.FileSize, len(entry.Content)), nil + } + + remainingInlineBudget := fs.option.SaveToFilerLimit - int64(len(entry.Content)) + if remainingInlineBudget > 0 { + // Read at most the remaining inline budget. Reaching the limit means + // the resulting file must be chunked, even if this is also EOF. + prefix, readErr := io.ReadAll(io.LimitReader(reader, remainingInlineBudget)) + if readErr != nil { + return nil, nil, 0, fmt.Errorf("read input: %w", readErr), nil + } + if int64(len(prefix)) < remainingInlineBudget { + inlineAppend := append([]byte{}, prefix...) + md5Hash = md5.New() + _, _ = md5Hash.Write(inlineAppend) + if len(inlineAppend) > 0 { + stats.FilerHandlerCounter.WithLabelValues(stats.ContentSaveToFiler).Inc() + } + return nil, md5Hash, int64(len(inlineAppend)), nil, inlineAppend + } + if chunkSize <= 0 { + return nil, nil, 0, fmt.Errorf("invalid chunk size %d for inline append promotion", chunkSize), nil + } + reader = io.MultiReader(bytes.NewReader(prefix), reader) + } else { + // Probe one byte so an empty append can preserve the inline entry + // even when the configured threshold is disabled or below its size. + prefix, readErr := io.ReadAll(io.LimitReader(reader, 1)) + if readErr != nil { + return nil, nil, 0, fmt.Errorf("read input: %w", readErr), nil + } + if len(prefix) == 0 { + return nil, md5.New(), 0, nil, []byte{} + } + if chunkSize <= 0 { + return nil, nil, 0, fmt.Errorf("invalid chunk size %d for inline append promotion", chunkSize), nil + } + reader = io.MultiReader(bytes.NewReader(prefix), reader) + } + } + } + return fs.uploadReaderToChunks(ctx, r, reader, chunkOffset, chunkSize, fileName, contentType, isAppend, so) }