seaweedfs/weed/filer/filer_notify_append.go

87 lines
2.4 KiB
Go
Raw Normal View History

2020-09-01 15:21:19 +08:00
package filer
2020-03-30 16:19:33 +08:00
import (
"context"
"fmt"
"os"
"time"
"github.com/seaweedfs/seaweedfs/weed/operation"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/util"
2020-03-30 16:19:33 +08:00
)
func (f *Filer) appendToFile(targetFile string, data []byte) error {
2020-11-16 08:58:48 +08:00
assignResult, uploadResult, err2 := f.assignAndUpload(targetFile, data)
2020-04-13 05:03:07 +08:00
if err2 != nil {
return err2
2020-03-30 16:19:33 +08:00
}
// find out existing entry
fullpath := util.FullPath(targetFile)
entry, err := f.FindEntry(context.Background(), fullpath)
var offset int64 = 0
if err == filer_pb.ErrNotFound {
entry = &Entry{
FullPath: fullpath,
Attr: Attr{
Crtime: time.Now(),
Mtime: time.Now(),
Mode: os.FileMode(0644),
Uid: OS_UID,
Gid: OS_GID,
},
}
2021-09-28 14:59:45 +08:00
} else if err != nil {
return fmt.Errorf("find %s: %v", fullpath, err)
2020-03-30 16:19:33 +08:00
} else {
offset = int64(TotalSize(entry.Chunks))
}
// append to existing chunks
entry.Chunks = append(entry.Chunks, uploadResult.ToPbFileChunk(assignResult.Fid, offset))
2020-03-30 16:19:33 +08:00
// update the entry
err = f.CreateEntry(context.Background(), entry, false, false, nil, false)
2020-03-30 16:19:33 +08:00
return err
}
2020-04-13 05:03:07 +08:00
2020-11-16 08:58:48 +08:00
func (f *Filer) assignAndUpload(targetFile string, data []byte) (*operation.AssignResult, *operation.UploadResult, error) {
2020-04-13 05:03:07 +08:00
// assign a volume location
2020-11-16 08:58:48 +08:00
rule := f.FilerConf.MatchStorageRule(targetFile)
2020-04-13 05:03:07 +08:00
assignRequest := &operation.VolumeAssignRequest{
Count: 1,
2020-11-16 08:58:48 +08:00
Collection: util.Nvl(f.metaLogCollection, rule.Collection),
Replication: util.Nvl(f.metaLogReplication, rule.Replication),
WritableVolumeCount: rule.VolumeGrowthCount,
2020-04-13 05:03:07 +08:00
}
2020-11-16 08:58:48 +08:00
assignResult, err := operation.Assign(f.GetMaster, f.GrpcDialOption, assignRequest)
2020-04-13 05:03:07 +08:00
if err != nil {
2020-04-17 15:00:48 +08:00
return nil, nil, fmt.Errorf("AssignVolume: %v", err)
2020-04-13 05:03:07 +08:00
}
if assignResult.Error != "" {
2020-04-17 15:00:48 +08:00
return nil, nil, fmt.Errorf("AssignVolume error: %v", assignResult.Error)
2020-04-13 05:03:07 +08:00
}
// upload data
targetUrl := "http://" + assignResult.Url + "/" + assignResult.Fid
2021-09-07 07:20:49 +08:00
uploadOption := &operation.UploadOption{
UploadUrl: targetUrl,
Filename: "",
Cipher: f.Cipher,
IsInputCompressed: false,
MimeType: "",
PairMap: nil,
Jwt: assignResult.Auth,
}
uploadResult, err := operation.UploadData(data, uploadOption)
2020-04-13 05:03:07 +08:00
if err != nil {
2020-04-17 15:00:48 +08:00
return nil, nil, fmt.Errorf("upload data %s: %v", targetUrl, err)
2020-04-13 05:03:07 +08:00
}
// println("uploaded to", targetUrl)
2020-04-17 15:00:48 +08:00
return assignResult, uploadResult, nil
2020-04-13 05:03:07 +08:00
}