mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-09 07:47:54 +02:00
feat(s3/lifecycle): emit MPU init records from FilerListFunc
Last gap in the filer-backed ListFunc. A directory at .uploads/<id> carrying ExtMultipartObjectKey is the MPU init record; emit one bootstrap.Entry with IsMPUInit=true and DestKey set to the user's intended path. The walker's MatchPath uses DestKey for prefix matching; the WalkerDispatcher uses it for the LifecycleDelete RPC's ObjectPath. .uploads/<id> directories without the extended key are mid-write before metadata landed and stay skipped. isMPUInitDir is upgraded from the path-shape-only stub to the full shape + extended-attr check that mirrors router.mpuInitInfo and scheduler/bootstrap.go's same-named helper. Tests pin: valid init record emits with the right DestKey, missing ExtMultipartObjectKey skips the directory.
This commit is contained in:
1 parent
2cc8a56f86
commit
e3c5f624b9
2 files changed
+68
-16
No files matched your search
@@ -28,9 +28,9 @@ func init() { listPageSize.Store(1024) }
|
||||
// IsLatest / NumVersions / NoncurrentIndex / SuccessorModTime so the
|
||||
// walker's NoncurrentDays evaluation has the same per-version state
|
||||
// the streaming bootstrap injects via reader.Event.BootstrapVersion.
|
||||
//
|
||||
// MPU init records at .uploads/<id> are skipped here; the follow-up
|
||||
// commit emits one Entry per init with IsMPUInit and DestKey set.
|
||||
// MPU init records at .uploads/<id> with ExtMultipartObjectKey set
|
||||
// are emitted whole with IsMPUInit=true and DestKey carrying the
|
||||
// user's intended path.
|
||||
func FilerListFunc(client filer_pb.SeaweedFilerClient, bucketsPath string) bootstrap.ListFunc {
|
||||
return func(ctx context.Context, bucket, start string, cb func(*bootstrap.Entry) error) error {
|
||||
if client == nil {
|
||||
@@ -78,9 +78,18 @@ func walkBucketTree(ctx context.Context, client filer_pb.SeaweedFilerClient, dir
|
||||
full := dir + "/" + e.Name
|
||||
key := strings.TrimPrefix(full, bucketRoot+"/")
|
||||
if e.IsDirectory {
|
||||
if isMPUInitDirShape(key) {
|
||||
// TODO(phase4b): emit MPU init as a single Entry.
|
||||
return nil
|
||||
if isMPUInitDir(key, e) {
|
||||
if start != "" && key <= start {
|
||||
return nil
|
||||
}
|
||||
destKey := string(e.Extended[s3_constants.ExtMultipartObjectKey])
|
||||
return cb(&bootstrap.Entry{
|
||||
Path: key,
|
||||
DestKey: destKey,
|
||||
IsMPUInit: true,
|
||||
ModTime: time.Unix(e.Attributes.Mtime, int64(e.Attributes.MtimeNs)),
|
||||
Size: int64(e.Attributes.FileSize),
|
||||
})
|
||||
}
|
||||
return walkBucketTree(ctx, client, full, bucketRoot, start, cb)
|
||||
}
|
||||
@@ -282,14 +291,20 @@ func isVersionsDir(entry *filer_pb.Entry) bool {
|
||||
return entry.IsDirectory && strings.HasSuffix(entry.Name, s3_constants.VersionsFolder)
|
||||
}
|
||||
|
||||
// isMPUInitDirShape mirrors scheduler/bootstrap.go's isMPUInitDir but
|
||||
// checks the path shape only; the Extended-attr verification lives in
|
||||
// the eventual full-MPU emission path.
|
||||
func isMPUInitDirShape(key string) bool {
|
||||
// isMPUInitDir mirrors router.mpuInitInfo: a directory at
|
||||
// .uploads/<id> carrying the destination key in Extended is the MPU
|
||||
// init record. Verified shape + presence of ExtMultipartObjectKey;
|
||||
// directories at .uploads/<id> without the key are mid-write before
|
||||
// metadata landed and stay out of the dispatch path.
|
||||
func isMPUInitDir(key string, entry *filer_pb.Entry) bool {
|
||||
uploadsPrefix := s3_constants.MultipartUploadsFolder + "/"
|
||||
if !strings.HasPrefix(key, uploadsPrefix) {
|
||||
return false
|
||||
}
|
||||
rest := key[len(uploadsPrefix):]
|
||||
return rest != "" && !strings.ContainsRune(rest, '/')
|
||||
if rest == "" || strings.ContainsRune(rest, '/') {
|
||||
return false
|
||||
}
|
||||
v, ok := entry.Extended[s3_constants.ExtMultipartObjectKey]
|
||||
return ok && len(v) > 0
|
||||
}
|
||||
@@ -143,17 +143,54 @@ func TestFilerListFunc_RecursesIntoSubdirs(t *testing.T) {
|
||||
assert.Equal(t, []string{"logs/2026/b.log", "logs/a.log", "root.txt"}, paths)
|
||||
}
|
||||
|
||||
func TestFilerListFunc_SkipsUploadsDirsForNow(t *testing.T) {
|
||||
// Phase 4b-pre: `.uploads/<id>/` MPU init records are skipped
|
||||
// here. The follow-up commit emits one IsMPUInit Entry per init.
|
||||
func TestFilerListFunc_MPUInitEmitsWithDestKey(t *testing.T) {
|
||||
// .uploads/<id>/ with ExtMultipartObjectKey is the MPU init record.
|
||||
// One Entry per init, IsMPUInit=true, DestKey = the user's path.
|
||||
mtime := time.Now()
|
||||
mpuDir := dir("upload-id-1")
|
||||
mpuDir.Extended = map[string][]byte{
|
||||
s3_constants.ExtMultipartObjectKey: []byte("user/path/object"),
|
||||
}
|
||||
client := &fakeFiler{tree: map[string][]*filer_pb.Entry{
|
||||
"/buckets/bkt": {
|
||||
file("regular.txt", mtime, 1),
|
||||
dir(s3_constants.MultipartUploadsFolder),
|
||||
},
|
||||
"/buckets/bkt/" + s3_constants.MultipartUploadsFolder: {
|
||||
dir("upload-id-1"),
|
||||
mpuDir,
|
||||
},
|
||||
}}
|
||||
listFn := FilerListFunc(client, "/buckets")
|
||||
var got []*bootstrap.Entry
|
||||
require.NoError(t, listFn(context.Background(), "bkt", "", func(e *bootstrap.Entry) error {
|
||||
got = append(got, e)
|
||||
return nil
|
||||
}))
|
||||
require.Len(t, got, 2)
|
||||
byPath := map[string]*bootstrap.Entry{}
|
||||
for _, e := range got {
|
||||
byPath[e.Path] = e
|
||||
}
|
||||
require.NotNil(t, byPath["regular.txt"])
|
||||
require.NotNil(t, byPath[s3_constants.MultipartUploadsFolder+"/upload-id-1"])
|
||||
mpu := byPath[s3_constants.MultipartUploadsFolder+"/upload-id-1"]
|
||||
assert.True(t, mpu.IsMPUInit)
|
||||
assert.Equal(t, "user/path/object", mpu.DestKey)
|
||||
}
|
||||
|
||||
func TestFilerListFunc_MPUInitWithoutDestKeyIsSkipped(t *testing.T) {
|
||||
// A `.uploads/<id>` directory missing ExtMultipartObjectKey is
|
||||
// mid-write before metadata landed; the dispatcher would error
|
||||
// on empty DestKey. Skip it so the walk doesn't halt.
|
||||
mtime := time.Now()
|
||||
mpuDir := dir("upload-id-2") // no Extended
|
||||
client := &fakeFiler{tree: map[string][]*filer_pb.Entry{
|
||||
"/buckets/bkt": {
|
||||
file("regular.txt", mtime, 1),
|
||||
dir(s3_constants.MultipartUploadsFolder),
|
||||
},
|
||||
"/buckets/bkt/" + s3_constants.MultipartUploadsFolder: {
|
||||
mpuDir,
|
||||
},
|
||||
}}
|
||||
listFn := FilerListFunc(client, "/buckets")
|
||||
@@ -162,7 +199,7 @@ func TestFilerListFunc_SkipsUploadsDirsForNow(t *testing.T) {
|
||||
paths = append(paths, e.Path)
|
||||
return nil
|
||||
}))
|
||||
assert.Equal(t, []string{"regular.txt"}, paths, ".uploads/ must not surface raw children")
|
||||
assert.Equal(t, []string{"regular.txt"}, paths)
|
||||
}
|
||||
|
||||
// fileWithExt is a versioned-file entry helper.
|
||||
|
||||
Reference in new issue
Block a user