diff --git a/cmd/client/main.go b/cmd/client/main.go index daad1adb..249e6e93 100644 --- a/cmd/client/main.go +++ b/cmd/client/main.go @@ -265,11 +265,11 @@ func listFunc(con net.Conn, instanceCnf *config.Instance, args []string) error { return nil } -// Request to delete a specific storage object -func sendDeleteChunkRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { +// Request to remove a storage object or collect garbage. +func sendCleanupRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { ylogger.Zero.Info().Msg("Execute delete command") ylogger.Zero.Info().Str("name", args[0]).Msg("delete") - dmsg := message.NewDeleteMessage(args[0], segmentPort, segmentNum, confirm, garbage) + dmsg := message.BuildCleanupMessage(args[0], segmentPort, segmentNum, confirm, garbage) dmsg.CrazyDrop = crazyDrop msg := dmsg.Encode() _, err := con.Write(msg) @@ -295,17 +295,17 @@ func sendDeleteChunkRequest(con net.Conn, instanceCnf *config.Instance, args []s return nil } -// Request to delete a set of trash objects by prefix -func sendDeleteTrashRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { - ylogger.Zero.Info().Msg("Execute delete2 command") - ylogger.Zero.Info().Str("name", args[0]).Msg("delete2") - msg := message.NewDelete2Message(args[0], confirm, garbage).Encode() +// Request to drop a set of trash objects by prefix +func sendDropRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { + ylogger.Zero.Info().Msg("Execute delete command") + ylogger.Zero.Info().Str("name", args[0]).Msg("delete") + msg := message.BuildDropMessage(args[0], confirm, garbage).Encode() _, err := con.Write(msg) // Send message with socket to server if err != nil { return err } - ylogger.Zero.Debug().Bytes("msg", msg).Msg("constructed delete2 msg") + ylogger.Zero.Debug().Bytes("msg", msg).Msg("constructed delete msg") client := client.NewYClient(con) protoReader := pio.NewProtoReader(client) @@ -317,7 +317,7 @@ func sendDeleteTrashRequest(con net.Conn, instanceCnf *config.Instance, args []s } if ansType != message.MessageTypeReadyForQuery { - return fmt.Errorf("failed to delete2, msg: %v", body) + return fmt.Errorf("failed to delete, msg: %v", body) } return nil @@ -403,8 +403,8 @@ var copyCmd = &cobra.Command{ var deleteCmd = &cobra.Command{ Use: "delete", - Short: "delete", - RunE: Runner(sendDeleteChunkRequest), + Short: "remove an object or collect garbage", + RunE: Runner(sendCleanupRequest), Args: cobra.ExactArgs(1), } @@ -434,10 +434,10 @@ var goolCmd = &cobra.Command{ RunE: Runner(goolFunc), } -var delete2Cmd = &cobra.Command{ - Use: "deleteTrash", - Short: "deleteTrash", - RunE: Runner(sendDeleteTrashRequest), +var deletePrefixCmd = &cobra.Command{ + Use: "deletePrefix", + Short: "physically delete objects by prefix", + RunE: Runner(sendDropRequest), Args: cobra.ExactArgs(1), // name_prefix } @@ -481,9 +481,9 @@ func init() { untrashifyCmd.PersistentFlags().BoolVarP(&confirm, "confirm", "", false, "confirm deletion") rootCmd.AddCommand(untrashifyCmd) - delete2Cmd.PersistentFlags().BoolVarP(&confirm, "confirm", "", false, "confirm deletion") - delete2Cmd.PersistentFlags().BoolVarP(&garbage, "garbage", "g", false, "delete garbage") - rootCmd.AddCommand(delete2Cmd) + deletePrefixCmd.PersistentFlags().BoolVarP(&confirm, "confirm", "", false, "confirm deletion") + deletePrefixCmd.PersistentFlags().BoolVarP(&garbage, "garbage", "g", false, "delete garbage") + rootCmd.AddCommand(deletePrefixCmd) } func main() { diff --git a/pkg/core/core.go b/pkg/core/core.go index 7c1bee59..b81a4468 100644 --- a/pkg/core/core.go +++ b/pkg/core/core.go @@ -220,7 +220,6 @@ func (instance *Instance) Run(instanceCnf *config.Instance) error { // Error return value of `instance.pool.Put` is not checked ylogger.Zero.Warn().Uint("id", ycl.ID()).Err(err).Msg("error putting client to pool") } - if err := proc.ProcConn(instance.ProtoMgr, s, bs, cr, ycl, &instanceCnf.VacuumCnf); err != nil { ylogger.Zero.Warn().Uint("id", ycl.ID()).Err(err).Msg("error serving client") } diff --git a/pkg/message/delete_message.go b/pkg/message/cleanup_message.go similarity index 70% rename from pkg/message/delete_message.go rename to pkg/message/cleanup_message.go index 26cc6ee9..edfeca56 100644 --- a/pkg/message/delete_message.go +++ b/pkg/message/cleanup_message.go @@ -4,7 +4,10 @@ import ( "encoding/binary" ) -type DeleteMessage struct { // Seg port +// CleanupMessage requests removal of an object or garbage collection. +// Garbage collection uses soft deletion by default; CrazyDrop switches it to +// hard deletion. A single-file request is handled by the storage deleter. +type CleanupMessage struct { // Seg port Name string // File path Port uint64 // Port segment/instance DB Segnum uint64 // Segment number @@ -13,10 +16,10 @@ type DeleteMessage struct { // Seg port CrazyDrop bool // For garbage mode: delete immediately instead of moving to trash } -var _ ProtoMessage = &DeleteMessage{} +var _ ProtoMessage = &CleanupMessage{} -func NewDeleteMessage(name string, port uint64, seg uint64, confirm bool, garbage bool) *DeleteMessage { - return &DeleteMessage{ +func BuildCleanupMessage(name string, port uint64, seg uint64, confirm bool, garbage bool) *CleanupMessage { + return &CleanupMessage{ Name: name, Port: port, Segnum: seg, @@ -25,7 +28,7 @@ func NewDeleteMessage(name string, port uint64, seg uint64, confirm bool, garbag } } -func (c *DeleteMessage) Encode() []byte { +func (c *CleanupMessage) Encode() []byte { bt := []byte{ byte(MessageTypeDelete), 0, @@ -60,7 +63,7 @@ func (c *DeleteMessage) Encode() []byte { return append(bs, bt...) } -func (c *DeleteMessage) Decode(body []byte) { +func (c *CleanupMessage) Decode(body []byte) { if body[1] == 1 { c.Confirm = true } diff --git a/pkg/message/delete2_message.go b/pkg/message/delete2_message.go deleted file mode 100644 index c938d030..00000000 --- a/pkg/message/delete2_message.go +++ /dev/null @@ -1,55 +0,0 @@ -package message - -import ( - "encoding/binary" -) - -type Delete2Message struct { //seg port - Prefix string - Confirm bool - Garbage bool -} - -var _ ProtoMessage = &Delete2Message{} - -func NewDelete2Message(prefix string, confirm bool, garbage bool) *Delete2Message { - return &Delete2Message{ - Prefix: prefix, - Confirm: confirm, - Garbage: garbage, - } -} - -func (c *Delete2Message) Encode() []byte { - bt := []byte{ - byte(MessageTypeDelete2), - 0, - 0, - 0, - } - - if c.Confirm { - bt[1] = 1 - } - if c.Garbage { - bt[2] = 1 - } - - bt = append(bt, []byte(c.Prefix)...) - bt = append(bt, 0) - - ln := len(bt) + 8 - bs := make([]byte, 8) - binary.BigEndian.PutUint64(bs, uint64(ln)) - return append(bs, bt...) -} - -func (c *Delete2Message) Decode(body []byte) { - if body[1] == 1 { - c.Confirm = true - } - if body[2] == 1 { - c.Garbage = true - } - c.Prefix, _ = GetCstring(body[4:]) -} diff --git a/pkg/message/drop_message.go b/pkg/message/drop_message.go new file mode 100644 index 00000000..52dbabe1 --- /dev/null +++ b/pkg/message/drop_message.go @@ -0,0 +1,58 @@ +package message + +import ( + "encoding/binary" +) + +// DropMessage requests physical deletion of objects under a prefix. +// It retains MessageTypeDelete2 on the wire for compatibility with existing +// clients; Delete2 is only the historical protocol name. +type DropMessage struct { // Seg port + Prefix string // Object key prefix to delete under + Confirm bool // Execute deletion; false means dry-run + Garbage bool // Restrict deletion to garbage (trash) objects past retention +} + +var _ ProtoMessage = &DropMessage{} + +func BuildDropMessage(prefix string, confirm bool, garbage bool) *DropMessage { + return &DropMessage{ + Prefix: prefix, + Confirm: confirm, + Garbage: garbage, + } +} + +func (c *DropMessage) Encode() []byte { + bt := []byte{ + byte(MessageTypeDelete2), + 0, + 0, + 0, + } + + if c.Confirm { + bt[1] = 1 + } + if c.Garbage { + bt[2] = 1 + } + + bt = append(bt, []byte(c.Prefix)...) + bt = append(bt, 0) + + ln := len(bt) + 8 + bs := make([]byte, 8) + binary.BigEndian.PutUint64(bs, uint64(ln)) + return append(bs, bt...) +} + +func (c *DropMessage) Decode(body []byte) { + if body[1] == 1 { + c.Confirm = true + } + if body[2] == 1 { + c.Garbage = true + } + c.Prefix, _ = GetCstring(body[4:]) +} diff --git a/pkg/message/message_test.go b/pkg/message/message_test.go index 588dada7..cac15551 100644 --- a/pkg/message/message_test.go +++ b/pkg/message/message_test.go @@ -434,12 +434,12 @@ func TestCopyMsg(t *testing.T) { func TestDeleteMsg(t *testing.T) { assert := assert.New(t) - msg := message.NewDeleteMessage("myname/mynextname", 5432, 42, true, true) + msg := message.BuildCleanupMessage("myname/mynextname", 5432, 42, true, true) body := msg.Encode() assert.Equal(body[8], byte(message.MessageTypeDelete)) - msg2 := message.DeleteMessage{} + msg2 := message.CleanupMessage{} msg2.Decode(body[8:]) assert.Equal("myname/mynextname", msg2.Name) @@ -449,15 +449,15 @@ func TestDeleteMsg(t *testing.T) { assert.True(msg2.Garbage) } -func TestDelete2Msg(t *testing.T) { +func TestDeleteMessage(t *testing.T) { assert := assert.New(t) - msg := message.NewDelete2Message("trash/prefix", true, false) + msg := message.BuildDropMessage("trash/prefix", true, false) body := msg.Encode() assert.Equal(body[8], byte(message.MessageTypeDelete2)) - msg2 := message.Delete2Message{} + msg2 := message.DropMessage{} msg2.Decode(body[8:]) assert.Equal("trash/prefix", msg2.Prefix) diff --git a/pkg/proc/delete_handler.go b/pkg/proc/delete_handler.go index 2d52507e..8d1471b9 100644 --- a/pkg/proc/delete_handler.go +++ b/pkg/proc/delete_handler.go @@ -20,8 +20,8 @@ import ( //go:generate mockgen -destination=../../../test/mocks/mock_object.go -package mocks -build_flags -mod=readonly github.com/wal-g/wal-g/pkg/storages/storage Object type GarbageMgr interface { - HandleDeleteGarbage(message.DeleteMessage) error - HandleDeleteFile(message.DeleteMessage) error + HandleGarbageCleanup(message.CleanupMessage) error + HandleFileDeletion(message.CleanupMessage) error HandleUntrashifyFile(message.UntrashifyMessage) error } @@ -104,16 +104,82 @@ func (dh *BasicGarbageMgr) HandleUntrashifyFile(msg message.UntrashifyMessage) e return nil } -/* - * The design looks an awkward, - * because the function depends on two different vacuum config: - * - the local dh.Cnf - * - the global config.InstanceConfig() - * Example: TestDeleteGarbageInBucketMovesObjectsWhenCrazyDropDisabled - */ -func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.DeleteMessage) error { +type garbageCleanupStrategy struct { + workerCount int + defaultWorkerCount int + failedActionMsg string + failedFilesMsg string + logCandidate func(bucket string, file *object.ObjectInfo) + execute func(bucket string, file *object.ObjectInfo) error + complete func(t *metrics.DeleteOpTracker, size int64) +} + +// Deletes garbage files immediately (hard delete, CrazyDrop). +func (dh *BasicGarbageMgr) hardDeleteStrategy() garbageCleanupStrategy { + s := dh.StorageInterractor + return garbageCleanupStrategy{ + workerCount: dh.Cnf.TrashDeleteWorkers, + defaultWorkerCount: config.DefaultTrashDeleteWorkers, + failedActionMsg: "failed to delete some files", + failedFilesMsg: "some files were not deleted", + logCandidate: func(bucket string, file *object.ObjectInfo) { + ylogger.Zero.Debug().Str("bucket", bucket).Str("file", file.Path).Msg("file will be deleted") + }, + execute: func(bucket string, file *object.ObjectInfo) error { + ylogger.Zero.Info(). + Str("bucket", bucket). + Str("path", file.Path). + Msg("immediately delete garbage file") + return s.DeleteObject(bucket, file.Path) + }, + complete: func(t *metrics.DeleteOpTracker, size int64) { + t.CompleteDeleted(size) + }, + } +} + +// Moves garbage files to the trash prefix (soft delete). +func (dh *BasicGarbageMgr) softDeleteStrategy(segnum int) garbageCleanupStrategy { + s := dh.StorageInterractor + return garbageCleanupStrategy{ + workerCount: dh.Cnf.TrashMoveWorkers, + defaultWorkerCount: config.DefaultTrashMoveWorkers, + failedActionMsg: "failed to move some files", + failedFilesMsg: "some files were not moved", + logCandidate: func(bucket string, file *object.ObjectInfo) { + ylogger.Zero.Debug(). + Str("bucket", bucket). + Str("file", file.Path). + Str("trash_path", TrashPathFromRegPath(file.Path, segnum)). + Msg("file will be moved to trash") + }, + execute: func(bucket string, file *object.ObjectInfo) error { + trashPath := TrashPathFromRegPath(file.Path, segnum) + ylogger.Zero.Debug(). + Str("bucket", bucket). + Str("path", file.Path). + Str("trash_path", trashPath). + Msg("move garbage file to trash") + return s.MoveObject(bucket, file.Path, trashPath) + }, + complete: func(t *metrics.DeleteOpTracker, size int64) { + t.CompleteMoved(size) + }, + } +} + +// chooseDeletionStrategy selects soft or hard deletion for this request. +func (dh *BasicGarbageMgr) chooseDeletionStrategy(msg message.CleanupMessage) garbageCleanupStrategy { + if msg.CrazyDrop { + return dh.hardDeleteStrategy() + } + return dh.softDeleteStrategy(int(msg.Segnum)) +} + +func (dh *BasicGarbageMgr) CleanupGarbageInBucket(bucket string, msg message.CleanupMessage) error { start := time.Now() t := metrics.NewDeleteOpTracker(bucket, "DELETE_GARBAGE") + strategy := dh.chooseDeletionStrategy(msg) fileList, err := dh.ListGarbageFiles(bucket, msg) if err != nil { @@ -126,7 +192,7 @@ func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.Dele ylogger.Zero.Info().Str("bucket", bucket).Int("files", len(fileList)).Int("uploads", len(uploads)).Msg("garbage delete started") for _, file := range fileList { - ylogger.Zero.Debug().Str("bucket", bucket).Bool("crazy mode", msg.CrazyDrop).Str("file", file.Path).Msg("file will be deleted") + strategy.logCandidate(bucket, file) } for _, upload := range uploads { ylogger.Zero.Info().Str("bucket", bucket).Str("uploadId", upload).Msg("upload will be aborted") @@ -137,65 +203,28 @@ func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.Dele return nil } - var ( - failedActionMsg string - failedFilesMsg string - workerCount int - defaultWorkerCount int - operate func(file *object.ObjectInfo) error - ) - - if msg.CrazyDrop { - failedActionMsg = "failed to delete some files" - failedFilesMsg = "some files were not deleted" - workerCount = dh.Cnf.TrashDeleteWorkers - defaultWorkerCount = config.DefaultTrashDeleteWorkers - - operate = func(file *object.ObjectInfo) error { - ylogger.Zero.Info(). - Str("bucket", bucket). - Str("path", file.Path). - Msg("immediately delete garbage file") - - opErr := dh.StorageInterractor.DeleteObject(bucket, file.Path) - if opErr == nil { - t.CompleteDeleted(file.Size) - } - - return opErr - } - } else { - failedActionMsg = "failed to move some files" - failedFilesMsg = "some files were not moved" - workerCount = dh.Cnf.TrashMoveWorkers - defaultWorkerCount = config.DefaultTrashMoveWorkers - - operate = func(file *object.ObjectInfo) error { - trashPath := TrashPathFromRegPath(file.Path, int(msg.Segnum)) - - ylogger.Zero.Debug(). - Str("bucket", bucket). - Str("path", file.Path). - Str("trash_path", trashPath). - Msg("move garbage file to trash") - - opErr := dh.StorageInterractor.MoveObject(bucket, file.Path, trashPath) - if opErr == nil { - t.CompleteMoved(file.Size) - } - - return opErr - } - } - deleted := 0 for retryCount := 0; len(fileList) > 0 && retryCount < 10; retryCount++ { batch := fileList - fileList, err = dh.garbageFilesParallel(bucket, batch, workerCount, defaultWorkerCount, operate, failedActionMsg, t) + fileList, err = dh.OperateFilesParallel( + bucket, + batch, + strategy.workerCount, + strategy.defaultWorkerCount, + func(file *object.ObjectInfo) error { + opErr := strategy.execute(bucket, file) + if opErr == nil { + strategy.complete(t, file.Size) + } + return opErr + }, + strategy.failedActionMsg, + t, + ) deleted += len(batch) - len(fileList) if err != nil { - ylogger.Zero.Error().Str("bucket", bucket).AnErr("err", err).Msg(failedActionMsg) + ylogger.Zero.Error().Str("bucket", bucket).AnErr("err", err).Msg(strategy.failedActionMsg) } } @@ -203,9 +232,9 @@ func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.Dele for _, file := range fileList { t.RecordFailed(file.Size) } - ylogger.Zero.Error().Str("bucket", bucket).Int("failed files count", len(fileList)).Msg(failedFilesMsg) - ylogger.Zero.Error().Str("bucket", bucket).Any("failed files", fileList).Msg(failedActionMsg) - return errors.Wrap(err, failedActionMsg) + ylogger.Zero.Error().Str("bucket", bucket).Int("failed files count", len(fileList)).Msg(strategy.failedFilesMsg) + ylogger.Zero.Error().Str("bucket", bucket).Any("failed files", fileList).Msg(strategy.failedActionMsg) + return errors.Wrap(err, strategy.failedActionMsg) } for key, uploadId := range uploads { @@ -219,15 +248,16 @@ func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.Dele return nil } -func (dh *BasicGarbageMgr) HandleDeleteGarbage(msg message.DeleteMessage) error { +func (dh *BasicGarbageMgr) HandleGarbageCleanup(msg message.CleanupMessage) error { for _, b := range dh.StorageInterractor.ListBuckets() { - if err := dh.DeleteGarbageInBucket(b, msg); err != nil { + if err := dh.CleanupGarbageInBucket(b, msg); err != nil { return err } } return nil } -func (dh *BasicGarbageMgr) ListDelete2Files(bucket string, msg message.Delete2Message) ([]*object.ObjectInfo, error) { + +func (dh *BasicGarbageMgr) ListDeletePrefixFiles(bucket string, msg message.DropMessage) ([]*object.ObjectInfo, error) { // Get first backup lsn var err error @@ -254,7 +284,7 @@ func (dh *BasicGarbageMgr) ListDelete2Files(bucket string, msg message.Delete2Me return filesToDelete, nil } -func (dh *BasicGarbageMgr) garbageFilesParallel( +func (dh *BasicGarbageMgr) OperateFilesParallel( bucket string, fileList []*object.ObjectInfo, workerCount int, @@ -316,8 +346,8 @@ func (dh *BasicGarbageMgr) garbageFilesParallel( return nil, nil } -func (dh *BasicGarbageMgr) garbageTrashParallel(bucket string, fileList []*object.ObjectInfo, t *metrics.DeleteOpTracker) ([]*object.ObjectInfo, error) { - return dh.garbageFilesParallel( +func (dh *BasicGarbageMgr) DeleteGarbageParallel(bucket string, fileList []*object.ObjectInfo, t *metrics.DeleteOpTracker) ([]*object.ObjectInfo, error) { + return dh.OperateFilesParallel( bucket, fileList, dh.Cnf.TrashDeleteWorkers, @@ -334,11 +364,11 @@ func (dh *BasicGarbageMgr) garbageTrashParallel(bucket string, fileList []*objec ) } -func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.Delete2Message) error { +func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.DropMessage) error { start := time.Now() t := metrics.NewDeleteOpTracker(bucket, "DELETE_PREFIX") - fileList, err := dh.ListDelete2Files(bucket, msg) // Return the list of files to be deleted + fileList, err := dh.ListDeletePrefixFiles(bucket, msg) // Return the list of files to be deleted if err != nil { return errors.Wrap(err, "failed to delete file") } @@ -379,7 +409,7 @@ func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.Delet toDelete := len(fileList) for retryCount := 0; len(fileList) > 0 && retryCount < 10; retryCount++ { - fileList, err = dh.garbageTrashParallel(bucket, fileList, t) + fileList, err = dh.DeleteGarbageParallel(bucket, fileList, t) if err != nil { ylogger.Zero.Error().Str("bucket", bucket).AnErr("err", err).Msg("failed to delete garbage file") } @@ -405,7 +435,7 @@ func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.Delet return nil } -func (dh *BasicGarbageMgr) HandleDelete2Prefix(msg message.Delete2Message) error { +func (dh *BasicGarbageMgr) HandleDeletePrefix(msg message.DropMessage) error { for _, b := range dh.StorageInterractor.ListBuckets() { if err := dh.DeletePrefixInBucket(b, msg); err != nil { return err @@ -413,7 +443,9 @@ func (dh *BasicGarbageMgr) HandleDelete2Prefix(msg message.Delete2Message) error } return nil } -func (dh *BasicGarbageMgr) HandleDeleteFile(msg message.DeleteMessage) error { + +// Delete a single external object by exact name. +func (dh *BasicGarbageMgr) HandleFileDeletion(msg message.CleanupMessage) error { if !msg.Confirm { return nil } @@ -427,7 +459,7 @@ func (dh *BasicGarbageMgr) HandleDeleteFile(msg message.DeleteMessage) error { return nil } -func (dh *BasicGarbageMgr) ListGarbageFiles(bucket string, msg message.DeleteMessage) ([]*object.ObjectInfo, error) { +func (dh *BasicGarbageMgr) ListGarbageFiles(bucket string, msg message.CleanupMessage) ([]*object.ObjectInfo, error) { procStartTime := time.Now() t := metrics.NewDeleteOpTracker(bucket, "DELETE_GARBAGE") diff --git a/pkg/proc/delete_handler_test.go b/pkg/proc/delete_handler_test.go index 61e700c2..b4e9ba14 100644 --- a/pkg/proc/delete_handler_test.go +++ b/pkg/proc/delete_handler_test.go @@ -44,7 +44,7 @@ func histogramSampleCount(t *testing.T, vec *prometheus.HistogramVec, labels pro func TestFilesToDeletion(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -95,7 +95,7 @@ func TestFilesToDeletion(t *testing.T) { func TestFilesToDeletionSkipsRecentlyCreatedFiles(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -143,7 +143,7 @@ func TestFilesToDeletionSkipsRecentlyCreatedFiles(t *testing.T) { func TestFilesToDeletionRespectsProtectionSecondsWindow(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -196,7 +196,7 @@ func TestFilesToDeletionRespectsProtectionSecondsWindow(t *testing.T) { func TestFilesToDeletionClampsNegativeProtectionSeconds(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -269,10 +269,10 @@ func TestTrashPathConversion(t *testing.T) { } } -func TestListDelete2Files(t *testing.T) { +func TestListDeletePrefixFiles(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -295,7 +295,7 @@ func TestListDelete2Files(t *testing.T) { Cnf: &config.Vacuum{CheckBackup: true}, } - actualFilesToDelete, err := handler.ListDelete2Files("", msg) + actualFilesToDelete, err := handler.ListDeletePrefixFiles("", msg) assert.NoError(t, err) assert.Equal(t, len(filesInStorage), len(actualFilesToDelete)) assert.Equal(t, filesInStorage, actualFilesToDelete) @@ -304,7 +304,7 @@ func TestListDelete2Files(t *testing.T) { func TestDeletePrefixInBucketDeletesAllFilesOnceInParallelGarbagePass(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -345,7 +345,7 @@ func TestDeletePrefixInBucketDeletesAllFilesOnceInParallelGarbagePass(t *testing func TestDeletePrefixInBucketRetriesFailedGarbageDeletes(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -388,7 +388,7 @@ func TestDeletePrefixInBucketRetriesFailedGarbageDeletes(t *testing.T) { func TestDeletePrefixInBucketReturnsFailedGarbageDeletesAfterRetries(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -422,7 +422,7 @@ func TestDeletePrefixInBucketReturnsFailedGarbageDeletesAfterRetries(t *testing. func TestDeletePrefixInBucketCapsWorkerCountToFileCount(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -460,7 +460,7 @@ func TestDeletePrefixInBucketCapsWorkerCountToFileCount(t *testing.T) { func TestDeletePrefixInBucketUsesDefaultWorkerCountWhenConfiguredZero(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -497,10 +497,10 @@ func TestDeletePrefixInBucketUsesDefaultWorkerCountWhenConfiguredZero(t *testing assert.Equal(t, 1, deleted["trash/c"]) } -func TestDeleteGarbageInBucketMovesObjectsWhenCrazyDropDisabled(t *testing.T) { +func TestCleanupGarbageInBucketMovesObjectsWhenCrazyDropDisabled(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -538,7 +538,7 @@ func TestDeleteGarbageInBucketMovesObjectsWhenCrazyDropDisabled(t *testing.T) { Cnf: &config.InstanceConfig().VacuumCnf, } - err := handler.DeleteGarbageInBucket("trash", msg) + err := handler.CleanupGarbageInBucket("trash", msg) assert.NoError(t, err) } @@ -658,12 +658,12 @@ func TestHandleUntrashifyFileLogsProgressEveryInterval(t *testing.T) { assert.Equal(t, 2, strings.Count(logBuf.String(), "untrashify progress")) } -func TestDeleteGarbageInBucketRecordsMetrics(t *testing.T) { +func TestCleanupGarbageInBucketRecordsMetrics(t *testing.T) { ctrl := gomock.NewController(t) const bucket = "garbage-metrics-bucket" - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -700,7 +700,7 @@ func TestDeleteGarbageInBucketRecordsMetrics(t *testing.T) { Cnf: &config.InstanceConfig().VacuumCnf, } - err := handler.DeleteGarbageInBucket(bucket, msg) + err := handler.CleanupGarbageInBucket(bucket, msg) assert.NoError(t, err) labels := prometheus.Labels{"bucket": bucket, "operation": "DELETE_GARBAGE"} @@ -726,7 +726,7 @@ func TestDeletePrefixInBucketRecordsMetrics(t *testing.T) { const bucket = "prefix-metrics-bucket" - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -776,7 +776,7 @@ func TestDeletePrefixInBucketRecordsKeptAfterRetriesExhausted(t *testing.T) { const bucket = "prefix-metrics-retry-bucket" - msg := message.Delete2Message{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -812,10 +812,10 @@ func TestDeletePrefixInBucketRecordsKeptAfterRetriesExhausted(t *testing.T) { assert.Equal(t, float64(1), testutil.ToFloat64(metrics.DeleteProcessKept.With(labels))) } -func TestDeleteGarbageInBucketRetriesFailedTrashMoves(t *testing.T) { +func TestCleanupGarbageInBucketRetriesFailedTrashMoves(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -869,7 +869,7 @@ func TestDeleteGarbageInBucketRetriesFailedTrashMoves(t *testing.T) { Cnf: &config.InstanceConfig().VacuumCnf, } - err := handler.DeleteGarbageInBucket("trash", msg) + err := handler.CleanupGarbageInBucket("trash", msg) assert.NoError(t, err) assert.Equal(t, 1, attempts[filesInStorage[0].Path]) assert.Equal(t, 2, attempts[filesInStorage[1].Path]) diff --git a/pkg/proc/interaction.go b/pkg/proc/interaction.go index 5fde9284..ed1bbccf 100644 --- a/pkg/proc/interaction.go +++ b/pkg/proc/interaction.go @@ -452,8 +452,8 @@ func (*ProtoMgrImpl) ProcessCopyExtended( return nil } -func (*ProtoMgrImpl) ProcessDeleteExtended( - msg message.DeleteMessage, +func (*ProtoMgrImpl) ProcessCleanupExtended( + msg message.CleanupMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, @@ -472,15 +472,18 @@ func (*ProtoMgrImpl) ProcessDeleteExtended( var ( logMsg string - handleDelete func(msg message.DeleteMessage) error + successMsg string + handleDelete func(msg message.CleanupMessage) error ) if msg.Garbage { logMsg = "requested to perform external storage VACUUM" - handleDelete = dh.HandleDeleteGarbage + successMsg = "Deleted garbage successfully" + handleDelete = dh.HandleGarbageCleanup } else { logMsg = "requested to remove external chunk" - handleDelete = dh.HandleDeleteFile + successMsg = "Deleted chunk successfully" + handleDelete = dh.HandleFileDeletion } ylogger.Zero.Debug(). @@ -501,17 +504,15 @@ func (*ProtoMgrImpl) ProcessDeleteExtended( } if !msg.Confirm { ylogger.Zero.Warn().Msg("It was a dry-run, nothing was deleted") - } else if msg.Garbage { - ylogger.Zero.Info().Msg("Deleted garbage successfully") } else { - ylogger.Zero.Info().Msg("Deleted chunk successfully") + ylogger.Zero.Info().Msg(successMsg) } return nil } -func (*ProtoMgrImpl) ProcessDelete2Extended( - msg message.Delete2Message, +func (*ProtoMgrImpl) ProcessDropExtended( + msg message.DropMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, @@ -537,7 +538,7 @@ func (*ProtoMgrImpl) ProcessDelete2Extended( Str("Name", msg.Prefix). Bool("confirm", msg.Confirm).Msg("requested to delete any files") } - err := dh.HandleDelete2Prefix(msg) + err := dh.HandleDeletePrefix(msg) if err != nil { _ = ycl.ReplyError(err, "failed to finish operation") return err @@ -853,16 +854,16 @@ func ProcConn( case message.MessageTypeDelete: // receive message - msg := message.DeleteMessage{} + msg := message.CleanupMessage{} msg.Decode(body) - err := m.ProcessDeleteExtended(msg, s, bs, ycl, cnf) + err := m.ProcessCleanupExtended(msg, s, bs, ycl, cnf) if err != nil { return err } case message.MessageTypeDelete2: - msg := message.Delete2Message{} + msg := message.DropMessage{} msg.Decode(body) - err := m.ProcessDelete2Extended(msg, s, bs, ycl, cnf) + err := m.ProcessDropExtended(msg, s, bs, ycl, cnf) if err != nil { return err } diff --git a/pkg/proto/mgr.go b/pkg/proto/mgr.go index c170f227..19b3586c 100644 --- a/pkg/proto/mgr.go +++ b/pkg/proto/mgr.go @@ -54,15 +54,15 @@ type ProtoMgr interface { ycl client.YproxyClient) error /* TODO: merge these two */ - ProcessDeleteExtended( - msg message.DeleteMessage, + ProcessCleanupExtended( + msg message.CleanupMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, cnf *config.Vacuum) error - ProcessDelete2Extended( - msg message.Delete2Message, + ProcessDropExtended( + msg message.DropMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, diff --git a/test/regress/tests/08_delete_trash_dryrun.sh b/test/regress/tests/08_delete_trash_dryrun.sh index a6bfb39a..e2392fa1 100755 --- a/test/regress/tests/08_delete_trash_dryrun.sh +++ b/test/regress/tests/08_delete_trash_dryrun.sh @@ -3,5 +3,5 @@ set -ex echo 'trash data' | yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal put 'trash/garbage_file' yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal list '' -yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal deleteTrash 'trash' +yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal deletePrefix 'trash' yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal list '' diff --git a/test/regress/tests/09_delete_trash_confirm.sh b/test/regress/tests/09_delete_trash_confirm.sh index 7bd7e8c8..6787bde9 100755 --- a/test/regress/tests/09_delete_trash_confirm.sh +++ b/test/regress/tests/09_delete_trash_confirm.sh @@ -3,5 +3,5 @@ set -ex echo 'trash data' | yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal put 'trash/garbage_file' yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal list '' -yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal deleteTrash 'trash' --confirm +yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal deletePrefix 'trash' --confirm yp-client --config test/regress/conf/yproxy_vacuum.yaml -l fatal list ''