You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

145 lines
3.6 KiB

6 years ago
6 years ago
6 years ago
6 years ago
  1. package operation
  2. import (
  3. "context"
  4. "errors"
  5. "fmt"
  6. "google.golang.org/grpc"
  7. "net/http"
  8. "strings"
  9. "sync"
  10. "time"
  11. "github.com/chrislusf/seaweedfs/weed/pb/volume_server_pb"
  12. )
  13. type DeleteResult struct {
  14. Fid string `json:"fid"`
  15. Size int `json:"size"`
  16. Status int `json:"status"`
  17. Error string `json:"error,omitempty"`
  18. }
  19. func ParseFileId(fid string) (vid string, key_cookie string, err error) {
  20. commaIndex := strings.Index(fid, ",")
  21. if commaIndex <= 0 {
  22. return "", "", errors.New("Wrong fid format.")
  23. }
  24. return fid[:commaIndex], fid[commaIndex+1:], nil
  25. }
  26. // DeleteFiles batch deletes a list of fileIds
  27. func DeleteFiles(master string, grpcDialOption grpc.DialOption, fileIds []string) ([]*volume_server_pb.DeleteResult, error) {
  28. lookupFunc := func(vids []string) (map[string]LookupResult, error) {
  29. return LookupVolumeIds(master, grpcDialOption, vids)
  30. }
  31. return DeleteFilesWithLookupVolumeId(grpcDialOption, fileIds, lookupFunc)
  32. }
  33. func DeleteFilesWithLookupVolumeId(grpcDialOption grpc.DialOption, fileIds []string, lookupFunc func(vid []string) (map[string]LookupResult, error)) ([]*volume_server_pb.DeleteResult, error) {
  34. var ret []*volume_server_pb.DeleteResult
  35. vid_to_fileIds := make(map[string][]string)
  36. var vids []string
  37. for _, fileId := range fileIds {
  38. vid, _, err := ParseFileId(fileId)
  39. if err != nil {
  40. ret = append(ret, &volume_server_pb.DeleteResult{
  41. FileId: fileId,
  42. Status: http.StatusBadRequest,
  43. Error: err.Error()},
  44. )
  45. continue
  46. }
  47. if _, ok := vid_to_fileIds[vid]; !ok {
  48. vid_to_fileIds[vid] = make([]string, 0)
  49. vids = append(vids, vid)
  50. }
  51. vid_to_fileIds[vid] = append(vid_to_fileIds[vid], fileId)
  52. }
  53. lookupResults, err := lookupFunc(vids)
  54. if err != nil {
  55. return ret, err
  56. }
  57. server_to_fileIds := make(map[string][]string)
  58. for vid, result := range lookupResults {
  59. if result.Error != "" {
  60. ret = append(ret, &volume_server_pb.DeleteResult{
  61. FileId: vid,
  62. Status: http.StatusBadRequest,
  63. Error: err.Error()},
  64. )
  65. continue
  66. }
  67. for _, location := range result.Locations {
  68. if _, ok := server_to_fileIds[location.Url]; !ok {
  69. server_to_fileIds[location.Url] = make([]string, 0)
  70. }
  71. server_to_fileIds[location.Url] = append(
  72. server_to_fileIds[location.Url], vid_to_fileIds[vid]...)
  73. }
  74. }
  75. var wg sync.WaitGroup
  76. for server, fidList := range server_to_fileIds {
  77. wg.Add(1)
  78. go func(server string, fidList []string) {
  79. defer wg.Done()
  80. if deleteResults, deleteErr := DeleteFilesAtOneVolumeServer(server, grpcDialOption, fidList); deleteErr != nil {
  81. err = deleteErr
  82. } else {
  83. ret = append(ret, deleteResults...)
  84. }
  85. }(server, fidList)
  86. }
  87. wg.Wait()
  88. return ret, err
  89. }
  90. // DeleteFilesAtOneVolumeServer deletes a list of files that is on one volume server via gRpc
  91. func DeleteFilesAtOneVolumeServer(volumeServer string, grpcDialOption grpc.DialOption, fileIds []string) (ret []*volume_server_pb.DeleteResult, err error) {
  92. err = WithVolumeServerClient(volumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
  93. ctx, cancel := context.WithTimeout(context.Background(), time.Duration(5*time.Second))
  94. defer cancel()
  95. req := &volume_server_pb.BatchDeleteRequest{
  96. FileIds: fileIds,
  97. }
  98. resp, err := volumeServerClient.BatchDelete(ctx, req)
  99. // fmt.Printf("deleted %v %v: %v\n", fileIds, err, resp)
  100. if err != nil {
  101. return err
  102. }
  103. ret = append(ret, resp.Results...)
  104. return nil
  105. })
  106. if err != nil {
  107. return
  108. }
  109. for _, result := range ret {
  110. if result.Error != "" && result.Error != "not found" {
  111. return nil, fmt.Errorf("delete fileId %s: %v", result.FileId, result.Error)
  112. }
  113. }
  114. return
  115. }