mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
fix(sdk): Trash Directory Structure Optimization. #1000528020
Signed-off-by: zhumingze <zhumingze@oppo.com>
This commit is contained in:
parent
8adaad9f90
commit
d250ed5b92
@ -266,7 +266,7 @@ func (m *metadataManager) HandleMetadataOperation(conn net.Conn, p *Packet, remo
|
||||
if err = m.metaNode.opLimiter.Wait(p.Opcode); err != nil {
|
||||
log.LogWarnf("action[HandleMetadataOperation] op rate limited, opCode[%v] remote[%v] err[%v]",
|
||||
p.Opcode, remoteAddr, err)
|
||||
p.PacketErrorWithBody(proto.OpTryOtherAddr, []byte("too many requests, please retry later"))
|
||||
p.PacketErrorWithBody(proto.OpLimitedIoErr, []byte("too many requests, please retry later"))
|
||||
m.respondToClientWithVer(conn, p)
|
||||
return err
|
||||
}
|
||||
|
||||
@ -13,7 +13,7 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
defaultOpLimitBurst = 1
|
||||
defaultOpLimitBurst = 512
|
||||
// Default rate limit for operations (IOPS)
|
||||
defaultAsyncOpLimit = 10000 // 10K IOPS for async operations
|
||||
)
|
||||
|
||||
@ -23,9 +23,11 @@ func TestOpLimiter_Basic(t *testing.T) {
|
||||
err = ol.SetLimiter(name, 1, 0)
|
||||
require.NoError(t, err)
|
||||
|
||||
// second call without tokens should be rate limited when timeout=0
|
||||
// consume burst tokens then expect rate limit
|
||||
code := proto.GOpInfo[name]
|
||||
require.NoError(t, ol.Wait(code))
|
||||
for i := 0; i < defaultOpLimitBurst; i++ {
|
||||
require.NoError(t, ol.Wait(code))
|
||||
}
|
||||
require.Error(t, ol.Wait(code))
|
||||
|
||||
// remove and ensure Wait passes (no limiter present)
|
||||
@ -40,9 +42,10 @@ func TestOpLimiter_TimeoutBranch(t *testing.T) {
|
||||
// 1 QPS, timeout 1s
|
||||
require.NoError(t, ol.SetLimiter(name, 1, 1))
|
||||
code := proto.GOpInfo[name]
|
||||
// consume one immediately
|
||||
require.NoError(t, ol.Wait(code))
|
||||
|
||||
// consume burst first
|
||||
for i := 0; i < defaultOpLimitBurst; i++ {
|
||||
require.NoError(t, ol.Wait(code))
|
||||
}
|
||||
start := time.Now()
|
||||
// next will wait up to 1s until token available, then succeed
|
||||
require.NoError(t, ol.Wait(code))
|
||||
|
||||
@ -770,8 +770,11 @@ func (mw *MetaWrapper) txDelete_ll(parentID uint64, name string, isDir bool, ful
|
||||
parentPathAbsolute := mw.getCurrentPath(parentID)
|
||||
err = mw.trashPolicy.MoveToTrash(parentPathAbsolute, parentID, name, isDir)
|
||||
if err != nil {
|
||||
if strings.Contains(err.Error(), "quota exceeded") || strings.Contains(err.Error(), "no space") {
|
||||
log.LogDebugf("Delete_ll: quota exceeded, delete %v directly, err %v", name, err.Error())
|
||||
// delete directly if trash is full (quota exceeded, no space, or inode full)
|
||||
if strings.Contains(err.Error(), syscall.EDQUOT.Error()) ||
|
||||
strings.Contains(err.Error(), syscall.ENOSPC.Error()) ||
|
||||
strings.Contains(err.Error(), syscall.ENOMEM.Error()) {
|
||||
log.LogDebugf("Delete_ll: trash full, delete %v directly, err %v", name, err.Error())
|
||||
goto deleteDirectly
|
||||
}
|
||||
log.LogErrorf("Delete_ll: MoveToTrash name %v failed %v", name, err)
|
||||
@ -911,9 +914,11 @@ func (mw *MetaWrapper) Delete_ll_EX(parentID uint64, name string, isDir bool, ve
|
||||
parentPathAbsolute := mw.getCurrentPath(parentID)
|
||||
err = mw.trashPolicy.MoveToTrash(parentPathAbsolute, parentID, name, isDir)
|
||||
if err != nil {
|
||||
// delete it directly if quota is exceeded
|
||||
if strings.Contains(err.Error(), "quota exceeded") || strings.Contains(err.Error(), "no space") {
|
||||
log.LogDebugf("Delete_ll: quota exceeded, delete %v directly, err %v", name, err.Error())
|
||||
// delete directly if trash is full (quota exceeded, no space, or inode full)
|
||||
if strings.Contains(err.Error(), syscall.EDQUOT.Error()) ||
|
||||
strings.Contains(err.Error(), syscall.ENOSPC.Error()) ||
|
||||
strings.Contains(err.Error(), syscall.ENOMEM.Error()) {
|
||||
log.LogDebugf("Delete_ll: trash full, delete %v directly, err %v", name, err.Error())
|
||||
goto deleteDirectly
|
||||
}
|
||||
log.LogErrorf("Delete_ll: MoveToTrash name %v failed %v", name, err)
|
||||
|
||||
@ -1,6 +1,8 @@
|
||||
package meta
|
||||
|
||||
import (
|
||||
"crypto/md5"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"os"
|
||||
@ -14,8 +16,6 @@ import (
|
||||
"github.com/cubefs/cubefs/proto"
|
||||
"github.com/google/uuid"
|
||||
|
||||
// "syscall"
|
||||
|
||||
"github.com/cubefs/cubefs/util/errors"
|
||||
"github.com/cubefs/cubefs/util/log"
|
||||
)
|
||||
@ -32,7 +32,10 @@ const (
|
||||
OriginalName = "OriginalName"
|
||||
DefaultReaddirLimit = 4096
|
||||
TrashPathIgnore = "trashPathIgnore"
|
||||
LockExpireSeconds = 3600 // 1 hour
|
||||
LockExpireSeconds = 3600 // 1 hour
|
||||
BucketRootPrefix = ".buckets" // bucket container under Current/Expired
|
||||
BucketHashWidth = 2 // first 2 hex chars
|
||||
RebuildDirName = "rebuild" // rebuilt result zone
|
||||
)
|
||||
|
||||
const (
|
||||
@ -91,6 +94,72 @@ func NewTrash(mw *MetaWrapper, interval int64, subDir string, traverseLimit int,
|
||||
return trash, nil
|
||||
}
|
||||
|
||||
// hashBucket returns the bucket name (hex) for a file, using parent path + file name to keep distribution stable.
|
||||
func (trash *Trash) hashBucket(parentPathAbsolute, fileName string) string {
|
||||
h := md5.Sum([]byte(path.Join(parentPathAbsolute, fileName)))
|
||||
return hex.EncodeToString(h[:])[:BucketHashWidth]
|
||||
}
|
||||
|
||||
// ensureChildDir creates a child directory under parentPath if absent, returning its full path and inode info.
|
||||
func (trash *Trash) ensureChildDir(parentPath, name string) (string, *proto.InodeInfo, error) {
|
||||
full := path.Join(parentPath, name)
|
||||
if info := trash.subDirCache.Get(full); info != nil {
|
||||
return full, info, nil
|
||||
}
|
||||
if trash.pathIsExist(full) {
|
||||
ino, err := trash.mw.LookupPath(full, true)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
info, err := trash.mw.InodeGet_ll(ino, true)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
trash.subDirCache.Put(full, info)
|
||||
return full, info, nil
|
||||
}
|
||||
|
||||
parentInfo, err := trash.LookupPath(parentPath, true)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
created, err := trash.CreateDirectory(parentInfo.Inode, name, parentInfo.Mode, parentInfo.Uid, parentInfo.Gid, full, true)
|
||||
if err != nil && err != syscall.EEXIST {
|
||||
return "", nil, err
|
||||
}
|
||||
if created == nil {
|
||||
ino, err := trash.mw.LookupPath(full, true)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
created, err = trash.mw.InodeGet_ll(ino, true)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
}
|
||||
trash.subDirCache.Put(full, created)
|
||||
return full, created, nil
|
||||
}
|
||||
|
||||
// ensureBucketDir prepares the bucket directory under Current (or provided root).
|
||||
func (trash *Trash) ensureBucketDir(bucketName string) (string, *proto.InodeInfo, error) {
|
||||
currentPath := path.Join(trash.trashRoot, CurrentName)
|
||||
if err := trash.createCurrent(true); err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
bucketRoot, _, err := trash.ensureChildDir(currentPath, BucketRootPrefix)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
return trash.ensureChildDir(bucketRoot, bucketName)
|
||||
}
|
||||
|
||||
// ensureRebuildRoot prepares the rebuild result root under the given base (Current or Expired_xxx).
|
||||
func (trash *Trash) ensureRebuildRoot(base string) (string, error) {
|
||||
_, _, err := trash.ensureChildDir(base, RebuildDirName)
|
||||
return path.Join(base, RebuildDirName), err
|
||||
}
|
||||
|
||||
func (trash *Trash) StartScheduleTask() {
|
||||
go trash.deleteWorker()
|
||||
go trash.buildDeletedFileParentDirsBackground()
|
||||
@ -201,11 +270,18 @@ retry:
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
goto retry
|
||||
}
|
||||
trashCurrentIno := trashCurrentInoInfo.Inode
|
||||
|
||||
bucketName := trash.hashBucket(parentPathAbsolute, fileName)
|
||||
bucketDir, bucketInfo, err := trash.ensureBucketDir(bucketName)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[MoveToTrash] ensureBucketDir failed: %v", err)
|
||||
return err
|
||||
}
|
||||
trashCurrentIno := bucketInfo.Inode
|
||||
srcPath := path.Join(trash.mountPath, parentPathAbsolute, fileName)
|
||||
// generate tmp file name
|
||||
tmpFileName := fmt.Sprintf("%v%v", trash.generateTmpFileName(parentPathAbsolute), fileName)
|
||||
dstPath := path.Join(trash.trashRoot, CurrentName, tmpFileName)
|
||||
dstPath := path.Join(bucketDir, tmpFileName)
|
||||
startCheck := time.Now()
|
||||
for {
|
||||
if trash.pathIsExistInTrash(dstPath) {
|
||||
@ -448,7 +524,7 @@ func (trash *Trash) deleteExpiredData() {
|
||||
continue
|
||||
}
|
||||
// extract timestamp from name
|
||||
err, checkPoint := trash.extractTimeStampFromName(entry.Name)
|
||||
checkPoint, err := trash.extractTimeStampFromName(entry.Name)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[deleteExpiredData]Extract timestamp from %s failed: %v", entry.Name, err.Error())
|
||||
continue
|
||||
@ -584,17 +660,17 @@ func (trash *Trash) removeAll(dirName string, dirIno uint64) {
|
||||
log.LogDebugf("action[deleteDir] delete complete %v", dirName)
|
||||
}
|
||||
|
||||
func (trash *Trash) extractTimeStampFromName(fileName string) (err error, timeStamp int64) {
|
||||
func (trash *Trash) extractTimeStampFromName(fileName string) (timeStamp int64, err error) {
|
||||
subs := strings.Split(fileName, "_")
|
||||
if len(subs) != 2 {
|
||||
return errors.New("fileName format is not valid"), 0
|
||||
return 0, errors.New("fileName format is not valid")
|
||||
}
|
||||
|
||||
parsedTime, err := time.ParseInLocation(ExpiredTimeFormat, subs[1], time.Local)
|
||||
if err != nil {
|
||||
return errors.New("fileName format is not valid"), 0
|
||||
return 0, errors.New("fileName format is not valid")
|
||||
}
|
||||
return nil, parsedTime.Unix()
|
||||
return parsedTime.Unix(), nil
|
||||
}
|
||||
|
||||
func (trash *Trash) pathIsExist(path string) bool {
|
||||
@ -764,8 +840,11 @@ func (trash *Trash) createParentPathInTrash(parentPath, rootDir string) (err err
|
||||
parentIno = info.Inode
|
||||
}
|
||||
continue
|
||||
} else if strings.Contains(err.Error(), "quota exceeded") || strings.Contains(err.Error(), "no space") {
|
||||
} else if strings.Contains(err.Error(), syscall.EDQUOT.Error()) ||
|
||||
strings.Contains(err.Error(), syscall.ENOSPC.Error()) ||
|
||||
strings.Contains(err.Error(), syscall.ENOMEM.Error()) {
|
||||
trash.deleteTask(parentIno, sub, true, path.Join(parentPath, sub))
|
||||
return
|
||||
} else {
|
||||
log.LogWarnf("action[createParentPathInTrash] CreateDirectory %v in trash failed: %v", cur, err.Error())
|
||||
return
|
||||
@ -880,6 +959,7 @@ type RebuildTask struct {
|
||||
Type uint32
|
||||
Inode uint64
|
||||
FileIno uint64
|
||||
SrcDir string
|
||||
}
|
||||
|
||||
func (trash *Trash) buildDeletedFileParentDirs() {
|
||||
@ -892,50 +972,30 @@ func (trash *Trash) buildDeletedFileParentDirs() {
|
||||
log.LogDebugf("action[buildDeletedFileParentDirs] start")
|
||||
trashCurrent := path.Join(trash.trashRoot, CurrentName)
|
||||
if !trash.pathIsExist(trashCurrent) {
|
||||
// log.LogWarnf("action[buildDeletedFileParentDirs] trashCurrent is not exist")
|
||||
return
|
||||
}
|
||||
// readdir
|
||||
var (
|
||||
trashInfo *proto.InodeInfo
|
||||
err error
|
||||
taskCh = make(chan RebuildTask, 1024)
|
||||
wg = sync.WaitGroup{}
|
||||
)
|
||||
if value := trash.subDirCache.Get(trashCurrent); value == nil {
|
||||
ino, _ := trash.mw.LookupPath(trashCurrent, true)
|
||||
trashInfo, err = trash.mw.InodeGet_ll(ino, true)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirs] get %v inode info failed:%v", trashCurrent, err.Error())
|
||||
return
|
||||
} else {
|
||||
trash.subDirCache.Put(trashCurrent, trashInfo)
|
||||
log.LogDebugf("action[buildDeletedFileParentDirs] store %v info %v", trashCurrent, trashInfo)
|
||||
}
|
||||
} else {
|
||||
trashInfo = value
|
||||
}
|
||||
if trashInfo == nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirs] trashInfo is nil %v", trashCurrent)
|
||||
|
||||
targetRoot, err := trash.ensureRebuildRoot(trashCurrent)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirs] ensureRebuildRoot failed: %v", err)
|
||||
return
|
||||
}
|
||||
var (
|
||||
noMore = false
|
||||
from = ""
|
||||
// children []proto.Dentry
|
||||
)
|
||||
// children, err := trash.mw.ReadDir_ll(trashInfo.Inode)
|
||||
// if err != nil {
|
||||
// log.LogWarnf("action[buildDeletedFileParentDirs] ReadDir %v failed:%v", trashCurrent, err.Error())
|
||||
// return
|
||||
// }
|
||||
|
||||
bucketRoot := path.Join(trashCurrent, BucketRootPrefix)
|
||||
if !trash.pathIsExist(bucketRoot) {
|
||||
log.LogDebugf("action[buildDeletedFileParentDirs] bucketRoot %v not exist", bucketRoot)
|
||||
return
|
||||
}
|
||||
|
||||
taskCh := make(chan RebuildTask, 1024)
|
||||
wg := sync.WaitGroup{}
|
||||
rebuildTaskFunc := func() {
|
||||
defer wg.Done()
|
||||
for task := range taskCh {
|
||||
if proto.IsDir(task.Type) {
|
||||
trash.rebuildDir(task.Name, trashCurrent, task.Inode, task.FileIno, false)
|
||||
trash.rebuildDir(task, targetRoot)
|
||||
} else {
|
||||
trash.rebuildFile(task.Name, trashCurrent, task.FileIno, false)
|
||||
trash.rebuildFile(task, targetRoot)
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -944,28 +1004,60 @@ func (trash *Trash) buildDeletedFileParentDirs() {
|
||||
wg.Add(1)
|
||||
go rebuildTaskFunc()
|
||||
}
|
||||
for !noMore {
|
||||
batches, err := trash.mw.ReadDirLimit_ll(trashInfo.Inode, from, DefaultReaddirLimit, true)
|
||||
if err != nil {
|
||||
log.LogErrorf("action[buildDeletedFileParentDirs] ReadDirLimit_ll: ino(%v) err(%v) from(%v)", trashInfo.Inode, err, from)
|
||||
return
|
||||
|
||||
// iterate buckets
|
||||
bucketInfo, err := trash.LookupPath(bucketRoot, true)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirs] lookup bucketRoot failed: %v", err)
|
||||
close(taskCh)
|
||||
wg.Wait()
|
||||
return
|
||||
}
|
||||
bucketEntries, err := trash.mw.ReadDir_ll(bucketInfo.Inode, true)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirs] read bucketRoot failed: %v", err)
|
||||
close(taskCh)
|
||||
wg.Wait()
|
||||
return
|
||||
}
|
||||
for _, bucket := range bucketEntries {
|
||||
if !proto.IsDir(bucket.Type) {
|
||||
continue
|
||||
}
|
||||
batchNr := uint64(len(batches))
|
||||
if batchNr == 0 || (from != "" && batchNr == 1) {
|
||||
break
|
||||
} else if batchNr < DefaultReaddirLimit {
|
||||
noMore = true
|
||||
}
|
||||
if from != "" {
|
||||
batches = batches[1:]
|
||||
}
|
||||
for _, child := range batches {
|
||||
log.LogDebugf("action[buildDeletedFileParentDirs] rebuild %v type %v", child.Name, child.Type)
|
||||
if strings.Contains(child.Name, ParentDirPrefix) || strings.Contains(child.Name, LongNamePrefix) {
|
||||
taskCh <- RebuildTask{Name: child.Name, Type: child.Type, Inode: trashInfo.Inode, FileIno: child.Inode}
|
||||
bucketPath := path.Join(bucketRoot, bucket.Name)
|
||||
var (
|
||||
noMore = false
|
||||
from = ""
|
||||
)
|
||||
for !noMore {
|
||||
batches, err := trash.mw.ReadDirLimit_ll(bucket.Inode, from, DefaultReaddirLimit, true)
|
||||
if err != nil {
|
||||
log.LogErrorf("action[buildDeletedFileParentDirs] ReadDirLimit_ll: bucket(%v) err(%v) from(%v)", bucketPath, err, from)
|
||||
break
|
||||
}
|
||||
batchNr := uint64(len(batches))
|
||||
if batchNr == 0 || (from != "" && batchNr == 1) {
|
||||
break
|
||||
} else if batchNr < DefaultReaddirLimit {
|
||||
noMore = true
|
||||
}
|
||||
if from != "" {
|
||||
batches = batches[1:]
|
||||
}
|
||||
for _, child := range batches {
|
||||
log.LogDebugf("action[buildDeletedFileParentDirs] rebuild %v type %v", child.Name, child.Type)
|
||||
if strings.Contains(child.Name, ParentDirPrefix) || strings.Contains(child.Name, LongNamePrefix) {
|
||||
taskCh <- RebuildTask{
|
||||
Name: child.Name,
|
||||
Type: child.Type,
|
||||
Inode: bucket.Inode,
|
||||
FileIno: child.Inode,
|
||||
SrcDir: bucketPath,
|
||||
}
|
||||
}
|
||||
}
|
||||
from = batches[len(batches)-1].Name
|
||||
}
|
||||
from = batches[len(batches)-1].Name
|
||||
}
|
||||
|
||||
close(taskCh)
|
||||
@ -991,7 +1083,17 @@ func (trash *Trash) buildDeletedFileParentDirsForExpired() {
|
||||
if !strings.HasPrefix(entry.Name, ExpiredPrefix) {
|
||||
continue
|
||||
}
|
||||
children, err := trash.mw.ReadDir_ll(entry.Inode, true)
|
||||
bucketRoot := path.Join(trash.trashRoot, entry.Name, BucketRootPrefix)
|
||||
if !trash.pathIsExist(bucketRoot) {
|
||||
continue
|
||||
}
|
||||
targetRoot, err := trash.ensureRebuildRoot(path.Join(trash.trashRoot, entry.Name))
|
||||
if err != nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirsForExpired] ensureRebuildRoot failed:%v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
bucketEntries, err := trash.mw.ReadDir_ll(entry.Inode, true)
|
||||
if err != nil {
|
||||
log.LogWarnf("action[buildDeletedFileParentDirsForExpired] ReadDir %v failed:%v", entry.Name, err.Error())
|
||||
continue
|
||||
@ -1004,9 +1106,9 @@ func (trash *Trash) buildDeletedFileParentDirsForExpired() {
|
||||
defer wg.Done()
|
||||
for task := range taskCh {
|
||||
if proto.IsDir(task.Type) {
|
||||
trash.rebuildDir(task.Name, path.Join(trash.trashRoot, entry.Name), task.Inode, task.FileIno, true)
|
||||
trash.rebuildDir(task, targetRoot)
|
||||
} else {
|
||||
trash.rebuildFile(task.Name, path.Join(trash.trashRoot, entry.Name), task.FileIno, true)
|
||||
trash.rebuildFile(task, targetRoot)
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -1014,10 +1116,37 @@ func (trash *Trash) buildDeletedFileParentDirsForExpired() {
|
||||
wg.Add(1)
|
||||
go rebuildTaskFunc()
|
||||
}
|
||||
for _, child := range children {
|
||||
log.LogDebugf("action[buildDeletedFileParentDirsForExpired] rebuild %v type %v", child.Name, child.Type)
|
||||
if strings.Contains(child.Name, ParentDirPrefix) || strings.Contains(child.Name, LongNamePrefix) {
|
||||
taskCh <- RebuildTask{Name: child.Name, Type: child.Type, Inode: entry.Inode, FileIno: child.Inode}
|
||||
for _, bucket := range bucketEntries {
|
||||
if !proto.IsDir(bucket.Type) || bucket.Name != BucketRootPrefix {
|
||||
continue
|
||||
}
|
||||
bucketPath := path.Join(trash.trashRoot, entry.Name, bucket.Name)
|
||||
var (
|
||||
noMore = false
|
||||
from = ""
|
||||
)
|
||||
for !noMore {
|
||||
batches, err := trash.mw.ReadDirLimit_ll(bucket.Inode, from, DefaultReaddirLimit, true)
|
||||
if err != nil {
|
||||
log.LogErrorf("action[buildDeletedFileParentDirsForExpired] ReadDirLimit_ll: bucket(%v) err(%v) from(%v)", bucketPath, err, from)
|
||||
break
|
||||
}
|
||||
batchNr := uint64(len(batches))
|
||||
if batchNr == 0 || (from != "" && batchNr == 1) {
|
||||
break
|
||||
} else if batchNr < DefaultReaddirLimit {
|
||||
noMore = true
|
||||
}
|
||||
if from != "" {
|
||||
batches = batches[1:]
|
||||
}
|
||||
for _, child := range batches {
|
||||
log.LogDebugf("action[buildDeletedFileParentDirsForExpired] rebuild %v type %v", child.Name, child.Type)
|
||||
if strings.Contains(child.Name, ParentDirPrefix) || strings.Contains(child.Name, LongNamePrefix) {
|
||||
taskCh <- RebuildTask{Name: child.Name, Type: child.Type, Inode: bucket.Inode, FileIno: child.Inode, SrcDir: bucketPath}
|
||||
}
|
||||
}
|
||||
from = batches[len(batches)-1].Name
|
||||
}
|
||||
}
|
||||
close(taskCh)
|
||||
@ -1026,64 +1155,50 @@ func (trash *Trash) buildDeletedFileParentDirsForExpired() {
|
||||
}
|
||||
}
|
||||
|
||||
func (trash *Trash) rebuildFile(fileName, trashCurrent string, fileIno uint64, forExpired bool) {
|
||||
log.LogDebugf("action[rebuildFile]: rebuild file %v in %v", fileName, trashCurrent)
|
||||
originName := fileName
|
||||
fileName = trash.recoverPosixPathName(fileName, fileIno)
|
||||
func (trash *Trash) rebuildFile(task RebuildTask, targetRoot string) {
|
||||
log.LogDebugf("action[rebuildFile]: rebuild file %v in %v", task.Name, task.SrcDir)
|
||||
originName := task.Name
|
||||
fileName := trash.recoverPosixPathName(task.Name, task.FileIno)
|
||||
log.LogDebugf("action[rebuildFile]: recover %v to %v", originName, fileName)
|
||||
parentDir := path.Dir(fileName)
|
||||
baseName := path.Base(fileName)
|
||||
|
||||
if parentDir == "." { // file in trash root
|
||||
if trash.pathIsExist(path.Join(trashCurrent, baseName)) {
|
||||
if trash.pathIsExist(path.Join(targetRoot, baseName)) {
|
||||
baseName = fmt.Sprintf("%s_%v", baseName, time.Now().Unix())
|
||||
}
|
||||
if err := trash.rename(path.Join(trashCurrent, originName), path.Join(trashCurrent, baseName)); err != nil {
|
||||
if err := trash.rename(path.Join(task.SrcDir, originName), path.Join(targetRoot, baseName)); err != nil {
|
||||
log.LogWarnf("action[rebuildFile]: recover %v to %v failed:err %v",
|
||||
path.Join(trashCurrent, originName), path.Join(trashCurrent, baseName), err)
|
||||
path.Join(task.SrcDir, originName), path.Join(targetRoot, baseName), err)
|
||||
}
|
||||
} else {
|
||||
//
|
||||
if forExpired {
|
||||
if !trash.pathIsExist(path.Join(trashCurrent, parentDir)) {
|
||||
// create in expired
|
||||
trash.createParentPathInTrash(parentDir, trashCurrent)
|
||||
log.LogDebugf("action[rebuildFile]: build parentDir %v in %v", parentDir, trashCurrent)
|
||||
}
|
||||
} else {
|
||||
if !trash.pathIsExistInTrash(parentDir) {
|
||||
trash.createParentPathInTrash(parentDir, CurrentName)
|
||||
}
|
||||
if !trash.pathIsExist(path.Join(targetRoot, parentDir)) {
|
||||
trash.createParentPathInTrash(parentDir, targetRoot)
|
||||
}
|
||||
if trash.pathIsExist(path.Join(trashCurrent, fileName)) {
|
||||
if trash.pathIsExist(path.Join(targetRoot, fileName)) {
|
||||
baseName = fmt.Sprintf("%s_%v", baseName, time.Now().Unix())
|
||||
}
|
||||
if err := trash.rename(path.Join(trashCurrent, originName), path.Join(trashCurrent, parentDir, baseName)); err != nil {
|
||||
if err := trash.rename(path.Join(task.SrcDir, originName), path.Join(targetRoot, parentDir, baseName)); err != nil {
|
||||
log.LogWarnf("action[rebuildFile]: recover %v to %v failed:err %v",
|
||||
path.Join(trashCurrent, originName), path.Join(trashCurrent, parentDir, baseName), err.Error())
|
||||
path.Join(task.SrcDir, originName), path.Join(targetRoot, parentDir, baseName), err.Error())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (trash *Trash) rebuildDir(dirName, trashCurrent string, ino uint64, fileIno uint64, forExpired bool) {
|
||||
log.LogDebugf("action[rebuildDir]: rebuild dir %v in %v", dirName, trashCurrent)
|
||||
originName := dirName
|
||||
dirName = trash.recoverPosixPathName(dirName, fileIno)
|
||||
func (trash *Trash) rebuildDir(task RebuildTask, targetRoot string) {
|
||||
log.LogDebugf("action[rebuildDir]: rebuild dir %v in %v", task.Name, task.SrcDir)
|
||||
originName := task.Name
|
||||
dirName := trash.recoverPosixPathName(task.Name, task.FileIno)
|
||||
var err error
|
||||
if forExpired {
|
||||
err = trash.createParentPathInTrash(dirName, trashCurrent)
|
||||
} else {
|
||||
err = trash.createParentPathInTrash(dirName, CurrentName)
|
||||
}
|
||||
err = trash.createParentPathInTrash(dirName, targetRoot)
|
||||
if err != nil {
|
||||
log.LogDebugf("action[rebuildDir]: createParentPathInTrash %v forExpired %v failed:err %v",
|
||||
dirName, forExpired, err)
|
||||
log.LogDebugf("action[rebuildDir]: createParentPathInTrash %v failed:err %v",
|
||||
dirName, err)
|
||||
return
|
||||
}
|
||||
log.LogDebugf("action[rebuildDir]: delete dir %v in %v[%v]", dirName, trashCurrent, ino)
|
||||
_, err = trash.mw.Delete_ll(ino, path.Base(originName), true, path.Join(trashCurrent, dirName), true)
|
||||
if err != nil {
|
||||
log.LogDebugf("action[rebuildDir]: delete dir %v in %v[%v] failed:err %v", dirName, trashCurrent, ino, err)
|
||||
log.LogDebugf("action[rebuildDir]: delete dir %v in %v[%v]", dirName, targetRoot, task.Inode)
|
||||
if _, err = trash.mw.Delete_ll(task.Inode, path.Base(originName), true, path.Join(task.SrcDir, originName), true); err != nil {
|
||||
log.LogDebugf("action[rebuildDir]: delete encoded dir %v failed: %v", path.Join(task.SrcDir, originName), err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
79
sdk/meta/trash_test.go
Normal file
79
sdk/meta/trash_test.go
Normal file
@ -0,0 +1,79 @@
|
||||
package meta
|
||||
|
||||
import (
|
||||
"path"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestHashBucket(t *testing.T) {
|
||||
trash := &Trash{}
|
||||
b1 := trash.hashBucket("/a/b", "c")
|
||||
b2 := trash.hashBucket("/a/b", "c")
|
||||
if b1 != b2 {
|
||||
t.Fatalf("hashBucket not deterministic: %s vs %s", b1, b2)
|
||||
}
|
||||
if len(b1) != BucketHashWidth {
|
||||
t.Fatalf("bucket length want %d got %d", BucketHashWidth, len(b1))
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecoverPosixPathNamePlain(t *testing.T) {
|
||||
trash := &Trash{}
|
||||
encoded := "a|__|b|__|c"
|
||||
got := trash.recoverPosixPathName(encoded, 0)
|
||||
if got != "a/b/c" {
|
||||
t.Fatalf("recoverPosixPathName want a/b/c got %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestGenerateTmpFileName(t *testing.T) {
|
||||
trash := &Trash{}
|
||||
|
||||
if got := trash.generateTmpFileName(""); got != ParentDirPrefix {
|
||||
t.Fatalf("root tmp name want %s got %s", ParentDirPrefix, got)
|
||||
}
|
||||
|
||||
// For nested path it should end with ParentDirPrefix and encode separators.
|
||||
got := trash.generateTmpFileName("a/b")
|
||||
if !strings.HasSuffix(got, ParentDirPrefix) {
|
||||
t.Fatalf("tmp name should end with ParentDirPrefix, got %s", got)
|
||||
}
|
||||
if !strings.Contains(got, ParentDirPrefix) {
|
||||
t.Fatalf("tmp name should contain encoded separator, got %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractTimeStampFromName(t *testing.T) {
|
||||
trash := &Trash{}
|
||||
now := time.Now().Unix()
|
||||
name := "Expired_" + time.Unix(now, 0).Format(ExpiredTimeFormat)
|
||||
ts, err := trash.extractTimeStampFromName(name)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected err: %v", err)
|
||||
}
|
||||
if ts != now {
|
||||
t.Fatalf("timestamp mismatch want %d got %d", now, ts)
|
||||
}
|
||||
|
||||
if _, err := trash.extractTimeStampFromName("bad_name"); err == nil {
|
||||
t.Fatalf("expect error for bad name")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTransferLongFileName(t *testing.T) {
|
||||
base := strings.Repeat("x", FileNameLengthMax+10)
|
||||
filePath := path.Join("/tmp", base)
|
||||
newName, oldName := transferLongFileName(filePath)
|
||||
|
||||
if oldName != base {
|
||||
t.Fatalf("old name want %s got %s", base, oldName)
|
||||
}
|
||||
if !strings.HasPrefix(newName, "/tmp/"+LongNamePrefix) {
|
||||
t.Fatalf("new name should start with long name prefix, got %s", newName)
|
||||
}
|
||||
if !strings.Contains(newName, ParentDirPrefix) {
|
||||
t.Fatalf("new name should contain ParentDirPrefix, got %s", newName)
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user