mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 06:22:05 +02:00
[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 <chrislusf@users.noreply.github.com> Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
3 files changed
+1340
-25
No files matched your search
@@ -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 {
|
||||
|
||||
File diff suppressed because it is too large.
Load diff
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user