seaweedfs/weed/operation/delete_content.go

150 lines
3.8 KiB
Go
Raw Normal View History

2012-09-26 18:27:10 +08:00
package operation
import (
2018-10-14 15:30:20 +08:00
"context"
2014-04-16 00:09:40 +08:00
"errors"
2018-10-14 15:30:20 +08:00
"fmt"
"github.com/chrislusf/seaweedfs/weed/glog"
"github.com/chrislusf/seaweedfs/weed/pb/volume_server_pb"
2019-02-19 04:11:52 +08:00
"google.golang.org/grpc"
2018-10-14 15:30:20 +08:00
"net/http"
2014-04-16 00:09:40 +08:00
"strings"
"sync"
2012-09-26 18:27:10 +08:00
)
2014-04-16 00:09:40 +08:00
type DeleteResult struct {
Fid string `json:"fid"`
Size int `json:"size"`
Status int `json:"status"`
Error string `json:"error,omitempty"`
2014-04-16 00:09:40 +08:00
}
func ParseFileId(fid string) (vid string, key_cookie string, err error) {
commaIndex := strings.Index(fid, ",")
if commaIndex <= 0 {
return "", "", errors.New("Wrong fid format.")
}
return fid[:commaIndex], fid[commaIndex+1:], nil
}
// DeleteFiles batch deletes a list of fileIds
2019-02-19 04:11:52 +08:00
func DeleteFiles(master string, grpcDialOption grpc.DialOption, fileIds []string) ([]*volume_server_pb.DeleteResult, error) {
2018-11-21 12:56:28 +08:00
lookupFunc := func(vids []string) (map[string]LookupResult, error) {
2019-02-19 04:11:52 +08:00
return LookupVolumeIds(master, grpcDialOption, vids)
2018-11-21 12:56:28 +08:00
}
2019-02-19 04:11:52 +08:00
return DeleteFilesWithLookupVolumeId(grpcDialOption, fileIds, lookupFunc)
2018-11-21 12:56:28 +08:00
}
2019-02-19 04:11:52 +08:00
func DeleteFilesWithLookupVolumeId(grpcDialOption grpc.DialOption, fileIds []string, lookupFunc func(vid []string) (map[string]LookupResult, error)) ([]*volume_server_pb.DeleteResult, error) {
2018-11-21 12:56:28 +08:00
var ret []*volume_server_pb.DeleteResult
2014-04-16 00:09:40 +08:00
vid_to_fileIds := make(map[string][]string)
var vids []string
for _, fileId := range fileIds {
vid, _, err := ParseFileId(fileId)
if err != nil {
ret = append(ret, &volume_server_pb.DeleteResult{
2019-02-15 16:09:19 +08:00
FileId: fileId,
Status: http.StatusBadRequest,
Error: err.Error()},
)
2014-04-16 00:09:40 +08:00
continue
}
if _, ok := vid_to_fileIds[vid]; !ok {
vid_to_fileIds[vid] = make([]string, 0)
vids = append(vids, vid)
}
vid_to_fileIds[vid] = append(vid_to_fileIds[vid], fileId)
}
2018-11-21 12:56:28 +08:00
lookupResults, err := lookupFunc(vids)
2014-04-16 00:09:40 +08:00
if err != nil {
return ret, err
}
server_to_fileIds := make(map[string][]string)
for vid, result := range lookupResults {
if result.Error != "" {
ret = append(ret, &volume_server_pb.DeleteResult{
FileId: vid,
Status: http.StatusBadRequest,
Error: err.Error()},
)
2014-04-16 00:09:40 +08:00
continue
}
for _, location := range result.Locations {
if _, ok := server_to_fileIds[location.Url]; !ok {
server_to_fileIds[location.Url] = make([]string, 0)
2014-04-16 00:09:40 +08:00
}
server_to_fileIds[location.Url] = append(
server_to_fileIds[location.Url], vid_to_fileIds[vid]...)
2014-04-16 00:09:40 +08:00
}
}
resultChan := make(chan []*volume_server_pb.DeleteResult, len(server_to_fileIds))
2014-04-16 00:09:40 +08:00
var wg sync.WaitGroup
for server, fidList := range server_to_fileIds {
wg.Add(1)
go func(server string, fidList []string) {
defer wg.Done()
2019-02-19 04:11:52 +08:00
if deleteResults, deleteErr := DeleteFilesAtOneVolumeServer(server, grpcDialOption, fidList); deleteErr != nil {
err = deleteErr
} else {
resultChan <- deleteResults
2014-04-16 00:09:40 +08:00
}
2014-04-16 00:09:40 +08:00
}(server, fidList)
}
wg.Wait()
close(resultChan)
for result := range resultChan {
ret = append(ret, result...)
}
glog.V(0).Infof("deleted %d items", len(ret))
2014-04-16 00:09:40 +08:00
return ret, err
}
// DeleteFilesAtOneVolumeServer deletes a list of files that is on one volume server via gRpc
2019-02-19 04:11:52 +08:00
func DeleteFilesAtOneVolumeServer(volumeServer string, grpcDialOption grpc.DialOption, fileIds []string) (ret []*volume_server_pb.DeleteResult, err error) {
2019-02-19 04:11:52 +08:00
err = WithVolumeServerClient(volumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
req := &volume_server_pb.BatchDeleteRequest{
FileIds: fileIds,
}
resp, err := volumeServerClient.BatchDelete(context.Background(), req)
2018-11-19 03:51:38 +08:00
// fmt.Printf("deleted %v %v: %v\n", fileIds, err, resp)
if err != nil {
return err
}
ret = append(ret, resp.Results...)
return nil
})
if err != nil {
return
}
for _, result := range ret {
2019-01-06 11:52:38 +08:00
if result.Error != "" && result.Error != "not found" {
return nil, fmt.Errorf("delete fileId %s: %v", result.FileId, result.Error)
}
}
return
2014-04-16 00:09:40 +08:00
}