From 32c7932348792b6ba5f722444908abd6eb794bd1 Mon Sep 17 00:00:00 2001 From: Vladislav Shchetinin Date: Mon, 17 Aug 2026 22:40:59 +0300 Subject: [PATCH 1/4] Refactor: Uniform appearance and semantics of function names and a simple strategy for choosing the type of deletion --- README.md | 2 +- cmd/client/main.go | 24 ++-- pkg/core/core.go | 1 - pkg/message/delete2_message.go | 19 +-- pkg/message/delete_message.go | 12 +- pkg/message/message_test.go | 8 +- pkg/proc/delete_handler.go | 120 ++++++++++++++---- pkg/proc/delete_handler_test.go | 40 +++--- pkg/proc/interaction.go | 29 +++-- pkg/proto/mgr.go | 8 +- test/regress/tests/08_delete_trash_dryrun.sh | 2 +- test/regress/tests/09_delete_trash_confirm.sh | 2 +- 12 files changed, 170 insertions(+), 97 deletions(-) diff --git a/README.md b/README.md index 711b6c56..024ee3de 100644 --- a/README.md +++ b/README.md @@ -110,7 +110,7 @@ DELETE OBSOLETE Delete/garbage-collection metrics are labeled per `bucket` (since yproxy can manage multiple storage buckets) and per `operation`, one of: ``` -DELETE_GARBAGE # DeleteGarbageInBucket +DELETE_GARBAGE # DisposalGarbageInBucket DELETE_PREFIX # DeletePrefixInBucket ``` `delete_request_latency_seconds` is additionally labeled by `stage` diff --git a/cmd/client/main.go b/cmd/client/main.go index daad1adb..a995859f 100644 --- a/cmd/client/main.go +++ b/cmd/client/main.go @@ -266,10 +266,10 @@ func listFunc(con net.Conn, instanceCnf *config.Instance, args []string) error { } // Request to delete a specific storage object -func sendDeleteChunkRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { +func sendDisposalRequest(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.BuildDisposalMessage(args[0], segmentPort, segmentNum, confirm, garbage) dmsg.CrazyDrop = crazyDrop msg := dmsg.Encode() _, err := con.Write(msg) @@ -296,16 +296,16 @@ func sendDeleteChunkRequest(con net.Conn, instanceCnf *config.Instance, args []s } // 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() +func sendDisposalPrefixRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { + ylogger.Zero.Info().Msg("Execute deletePrefix command") + ylogger.Zero.Info().Str("name", args[0]).Msg("deletePrefix") + msg := message.BuildDisposalPrefixMessage(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 deletePrefix 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 deletePrefix, msg: %v", body) } return nil @@ -404,7 +404,7 @@ var copyCmd = &cobra.Command{ var deleteCmd = &cobra.Command{ Use: "delete", Short: "delete", - RunE: Runner(sendDeleteChunkRequest), + RunE: Runner(sendDisposalRequest), Args: cobra.ExactArgs(1), } @@ -435,9 +435,9 @@ var goolCmd = &cobra.Command{ } var delete2Cmd = &cobra.Command{ - Use: "deleteTrash", - Short: "deleteTrash", - RunE: Runner(sendDeleteTrashRequest), + Use: "deletePrefix", + Short: "deletePrefix", + RunE: Runner(sendDisposalPrefixRequest), Args: cobra.ExactArgs(1), // name_prefix } 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/delete2_message.go b/pkg/message/delete2_message.go index c938d030..85066a7b 100644 --- a/pkg/message/delete2_message.go +++ b/pkg/message/delete2_message.go @@ -4,23 +4,24 @@ import ( "encoding/binary" ) -type Delete2Message struct { //seg port - Prefix string - Confirm bool - Garbage bool +// Requests deletion of all objects under a given prefix. +type DisposalPrefixMessage 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 = &Delete2Message{} +var _ ProtoMessage = &DisposalPrefixMessage{} -func NewDelete2Message(prefix string, confirm bool, garbage bool) *Delete2Message { - return &Delete2Message{ +func BuildDisposalPrefixMessage(prefix string, confirm bool, garbage bool) *DisposalPrefixMessage { + return &DisposalPrefixMessage{ Prefix: prefix, Confirm: confirm, Garbage: garbage, } } -func (c *Delete2Message) Encode() []byte { +func (c *DisposalPrefixMessage) Encode() []byte { bt := []byte{ byte(MessageTypeDelete2), 0, @@ -44,7 +45,7 @@ func (c *Delete2Message) Encode() []byte { return append(bs, bt...) } -func (c *Delete2Message) Decode(body []byte) { +func (c *DisposalPrefixMessage) Decode(body []byte) { if body[1] == 1 { c.Confirm = true } diff --git a/pkg/message/delete_message.go b/pkg/message/delete_message.go index 26cc6ee9..105f9efc 100644 --- a/pkg/message/delete_message.go +++ b/pkg/message/delete_message.go @@ -4,7 +4,7 @@ import ( "encoding/binary" ) -type DeleteMessage struct { // Seg port +type DisposalMessage struct { // Seg port Name string // File path Port uint64 // Port segment/instance DB Segnum uint64 // Segment number @@ -13,10 +13,10 @@ type DeleteMessage struct { // Seg port CrazyDrop bool // For garbage mode: delete immediately instead of moving to trash } -var _ ProtoMessage = &DeleteMessage{} +var _ ProtoMessage = &DisposalMessage{} -func NewDeleteMessage(name string, port uint64, seg uint64, confirm bool, garbage bool) *DeleteMessage { - return &DeleteMessage{ +func BuildDisposalMessage(name string, port uint64, seg uint64, confirm bool, garbage bool) *DisposalMessage { + return &DisposalMessage{ Name: name, Port: port, Segnum: seg, @@ -25,7 +25,7 @@ func NewDeleteMessage(name string, port uint64, seg uint64, confirm bool, garbag } } -func (c *DeleteMessage) Encode() []byte { +func (c *DisposalMessage) Encode() []byte { bt := []byte{ byte(MessageTypeDelete), 0, @@ -60,7 +60,7 @@ func (c *DeleteMessage) Encode() []byte { return append(bs, bt...) } -func (c *DeleteMessage) Decode(body []byte) { +func (c *DisposalMessage) Decode(body []byte) { if body[1] == 1 { c.Confirm = true } diff --git a/pkg/message/message_test.go b/pkg/message/message_test.go index 588dada7..2983c1fb 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.BuildDisposalMessage("myname/mynextname", 5432, 42, true, true) body := msg.Encode() assert.Equal(body[8], byte(message.MessageTypeDelete)) - msg2 := message.DeleteMessage{} + msg2 := message.DisposalMessage{} msg2.Decode(body[8:]) assert.Equal("myname/mynextname", msg2.Name) @@ -452,12 +452,12 @@ func TestDeleteMsg(t *testing.T) { func TestDelete2Msg(t *testing.T) { assert := assert.New(t) - msg := message.NewDelete2Message("trash/prefix", true, false) + msg := message.BuildDisposalPrefixMessage("trash/prefix", true, false) body := msg.Encode() assert.Equal(body[8], byte(message.MessageTypeDelete2)) - msg2 := message.Delete2Message{} + msg2 := message.DisposalPrefixMessage{} 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..55eb6dde 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 + HandleDisposalGarbage(message.DisposalMessage) error + HandleDeleteFile(message.DisposalMessage) error HandleUntrashifyFile(message.UntrashifyMessage) error } @@ -104,16 +104,75 @@ 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 fileDisposal struct { + workerCount int + defaultWorkerCount int + failedActionMsg string + failedFilesMsg string + logCandidate func(bucket string, file *object.ObjectInfo) + execute func(bucket string, file *object.ObjectInfo) error +} + +// Deletes garbage files immediately (hard delete, CrazyDrop). +func (dh *BasicGarbageMgr) deleteDisposal() fileDisposal { + s := dh.StorageInterractor + return fileDisposal{ + 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) + }, + } +} + +// Moves garbage files to the trash prefix (soft delete). +func (dh *BasicGarbageMgr) moveDisposal(segnum int) fileDisposal { + s := dh.StorageInterractor + return fileDisposal{ + 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) + }, + } +} + +// chooseDisposal picks how garbage files should be disposed of for this request. +func (dh *BasicGarbageMgr) chooseDisposal(msg message.DisposalMessage) fileDisposal { + if msg.CrazyDrop { + return dh.deleteDisposal() + } + return dh.moveDisposal(int(msg.Segnum)) +} + +func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.DisposalMessage) error { start := time.Now() t := metrics.NewDeleteOpTracker(bucket, "DELETE_GARBAGE") + disposal := dh.chooseDisposal(msg) fileList, err := dh.ListGarbageFiles(bucket, msg) if err != nil { @@ -126,7 +185,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") + disposal.logCandidate(bucket, file) } for _, upload := range uploads { ylogger.Zero.Info().Str("bucket", bucket).Str("uploadId", upload).Msg("upload will be aborted") @@ -192,10 +251,20 @@ func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.Dele 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, + disposal.workerCount, + disposal.defaultWorkerCount, + func(file *object.ObjectInfo) error { + return disposal.execute(bucket, file) + }, + disposal.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(disposal.failedActionMsg) } } @@ -219,15 +288,16 @@ func (dh *BasicGarbageMgr) DeleteGarbageInBucket(bucket string, msg message.Dele return nil } -func (dh *BasicGarbageMgr) HandleDeleteGarbage(msg message.DeleteMessage) error { +func (dh *BasicGarbageMgr) HandleDisposalGarbage(msg message.DisposalMessage) error { for _, b := range dh.StorageInterractor.ListBuckets() { - if err := dh.DeleteGarbageInBucket(b, msg); err != nil { + if err := dh.DisposalGarbageInBucket(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.DisposalPrefixMessage) ([]*object.ObjectInfo, error) { // Get first backup lsn var err error @@ -254,7 +324,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 +386,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 +404,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.DisposalPrefixMessage) 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") } @@ -405,7 +475,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.DisposalPrefixMessage) error { for _, b := range dh.StorageInterractor.ListBuckets() { if err := dh.DeletePrefixInBucket(b, msg); err != nil { return err @@ -413,7 +483,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) HandleDeleteFile(msg message.DisposalMessage) error { if !msg.Confirm { return nil } @@ -427,7 +499,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.DisposalMessage) ([]*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..24c4c451 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.DisposalMessage{ 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.DisposalMessage{ 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.DisposalMessage{ 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.DisposalMessage{ 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.DisposalPrefixMessage{ 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.DisposalPrefixMessage{ 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.DisposalPrefixMessage{ 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.DisposalPrefixMessage{ 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.DisposalPrefixMessage{ 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.DisposalPrefixMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -500,7 +500,7 @@ func TestDeletePrefixInBucketUsesDefaultWorkerCountWhenConfiguredZero(t *testing func TestDeleteGarbageInBucketMovesObjectsWhenCrazyDropDisabled(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.DisposalMessage{ 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.DisposalGarbageInBucket("trash", msg) assert.NoError(t, err) } @@ -663,7 +663,7 @@ func TestDeleteGarbageInBucketRecordsMetrics(t *testing.T) { const bucket = "garbage-metrics-bucket" - msg := message.DeleteMessage{ + msg := message.DisposalMessage{ 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.DisposalGarbageInBucket(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.DisposalPrefixMessage{ 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.DisposalPrefixMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -815,7 +815,7 @@ func TestDeletePrefixInBucketRecordsKeptAfterRetriesExhausted(t *testing.T) { func TestDeleteGarbageInBucketRetriesFailedTrashMoves(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DeleteMessage{ + msg := message.DisposalMessage{ 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.DisposalGarbageInBucket("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..57442678 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) ProcessDisposalExtended( + msg message.DisposalMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, @@ -472,14 +472,17 @@ func (*ProtoMgrImpl) ProcessDeleteExtended( var ( logMsg string - handleDelete func(msg message.DeleteMessage) error + successMsg string + handleDelete func(msg message.DisposalMessage) error ) if msg.Garbage { logMsg = "requested to perform external storage VACUUM" - handleDelete = dh.HandleDeleteGarbage + successMsg = "Deleted garbage successfully" + handleDelete = dh.HandleDisposalGarbage } else { logMsg = "requested to remove external chunk" + successMsg = "Deleted chunk successfully" handleDelete = dh.HandleDeleteFile } @@ -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) ProcessDisposalPrefixExtended( + msg message.DisposalPrefixMessage, 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.DisposalMessage{} msg.Decode(body) - err := m.ProcessDeleteExtended(msg, s, bs, ycl, cnf) + err := m.ProcessDisposalExtended(msg, s, bs, ycl, cnf) if err != nil { return err } case message.MessageTypeDelete2: - msg := message.Delete2Message{} + msg := message.DisposalPrefixMessage{} msg.Decode(body) - err := m.ProcessDelete2Extended(msg, s, bs, ycl, cnf) + err := m.ProcessDisposalPrefixExtended(msg, s, bs, ycl, cnf) if err != nil { return err } diff --git a/pkg/proto/mgr.go b/pkg/proto/mgr.go index c170f227..3221130e 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, + ProcessDisposalExtended( + msg message.DisposalMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, cnf *config.Vacuum) error - ProcessDelete2Extended( - msg message.Delete2Message, + ProcessDisposalPrefixExtended( + msg message.DisposalPrefixMessage, 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 '' From 46cba88c90798f074a82a0c72e9cadd82aa9c41c Mon Sep 17 00:00:00 2001 From: Vlasdislav Date: Fri, 4 Sep 2026 16:34:20 +0300 Subject: [PATCH 2/4] Refactor: Add the complete state to fileDisposal and delete rebase artifacts --- pkg/proc/delete_handler.go | 72 +++++++++----------------------------- 1 file changed, 16 insertions(+), 56 deletions(-) diff --git a/pkg/proc/delete_handler.go b/pkg/proc/delete_handler.go index 55eb6dde..153f783d 100644 --- a/pkg/proc/delete_handler.go +++ b/pkg/proc/delete_handler.go @@ -111,6 +111,7 @@ type fileDisposal struct { 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). @@ -131,6 +132,9 @@ func (dh *BasicGarbageMgr) deleteDisposal() fileDisposal { Msg("immediately delete garbage file") return s.DeleteObject(bucket, file.Path) }, + complete: func(t *metrics.DeleteOpTracker, size int64) { + t.CompleteDeleted(size) + }, } } @@ -158,6 +162,9 @@ func (dh *BasicGarbageMgr) moveDisposal(segnum int) fileDisposal { Msg("move garbage file to trash") return s.MoveObject(bucket, file.Path, trashPath) }, + complete: func(t *metrics.DeleteOpTracker, size int64) { + t.CompleteMoved(size) + }, } } @@ -196,57 +203,6 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di 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++ { @@ -257,7 +213,11 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di disposal.workerCount, disposal.defaultWorkerCount, func(file *object.ObjectInfo) error { - return disposal.execute(bucket, file) + opErr := disposal.execute(bucket, file) + if opErr == nil { + disposal.complete(t, file.Size) + } + return opErr }, disposal.failedActionMsg, t, @@ -272,9 +232,9 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di 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(disposal.failedFilesMsg) + ylogger.Zero.Error().Str("bucket", bucket).Any("failed files", fileList).Msg(disposal.failedActionMsg) + return errors.Wrap(err, disposal.failedActionMsg) } for key, uploadId := range uploads { @@ -449,7 +409,7 @@ func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.Dispo 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") } From ec75e92884c6e5438240f6de53da4d09b203b6de Mon Sep 17 00:00:00 2001 From: Vlasdislav Date: Fri, 4 Sep 2026 21:39:02 +0300 Subject: [PATCH 3/4] Refactor: Uniform appearance and semantics of function names --- README.md | 2 +- cmd/client/main.go | 36 +++++------ .../{delete_message.go => cleanup_message.go} | 15 +++-- .../{delete2_message.go => drop_message.go} | 16 ++--- pkg/message/message_test.go | 10 ++-- pkg/proc/delete_handler.go | 60 +++++++++---------- pkg/proc/delete_handler_test.go | 42 ++++++------- pkg/proc/interaction.go | 22 +++---- pkg/proto/mgr.go | 8 +-- 9 files changed, 108 insertions(+), 103 deletions(-) rename pkg/message/{delete_message.go => cleanup_message.go} (70%) rename pkg/message/{delete2_message.go => drop_message.go} (60%) diff --git a/README.md b/README.md index 024ee3de..711b6c56 100644 --- a/README.md +++ b/README.md @@ -110,7 +110,7 @@ DELETE OBSOLETE Delete/garbage-collection metrics are labeled per `bucket` (since yproxy can manage multiple storage buckets) and per `operation`, one of: ``` -DELETE_GARBAGE # DisposalGarbageInBucket +DELETE_GARBAGE # DeleteGarbageInBucket DELETE_PREFIX # DeletePrefixInBucket ``` `delete_request_latency_seconds` is additionally labeled by `stage` diff --git a/cmd/client/main.go b/cmd/client/main.go index a995859f..be91dd0a 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 sendDisposalRequest(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.BuildDisposalMessage(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 sendDisposalRequest(con net.Conn, instanceCnf *config.Instance, args []stri return nil } -// Request to delete a set of trash objects by prefix -func sendDisposalPrefixRequest(con net.Conn, instanceCnf *config.Instance, args []string) error { - ylogger.Zero.Info().Msg("Execute deletePrefix command") - ylogger.Zero.Info().Str("name", args[0]).Msg("deletePrefix") - msg := message.BuildDisposalPrefixMessage(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 deletePrefix msg") + ylogger.Zero.Debug().Bytes("msg", msg).Msg("constructed delete msg") client := client.NewYClient(con) protoReader := pio.NewProtoReader(client) @@ -317,7 +317,7 @@ func sendDisposalPrefixRequest(con net.Conn, instanceCnf *config.Instance, args } if ansType != message.MessageTypeReadyForQuery { - return fmt.Errorf("failed to deletePrefix, 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(sendDisposalRequest), + 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{ +var deletePrefixCmd = &cobra.Command{ Use: "deletePrefix", - Short: "deletePrefix", - RunE: Runner(sendDisposalPrefixRequest), + 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/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 105f9efc..edfeca56 100644 --- a/pkg/message/delete_message.go +++ b/pkg/message/cleanup_message.go @@ -4,7 +4,10 @@ import ( "encoding/binary" ) -type DisposalMessage 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 DisposalMessage struct { // Seg port CrazyDrop bool // For garbage mode: delete immediately instead of moving to trash } -var _ ProtoMessage = &DisposalMessage{} +var _ ProtoMessage = &CleanupMessage{} -func BuildDisposalMessage(name string, port uint64, seg uint64, confirm bool, garbage bool) *DisposalMessage { - return &DisposalMessage{ +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 BuildDisposalMessage(name string, port uint64, seg uint64, confirm bool, ga } } -func (c *DisposalMessage) Encode() []byte { +func (c *CleanupMessage) Encode() []byte { bt := []byte{ byte(MessageTypeDelete), 0, @@ -60,7 +63,7 @@ func (c *DisposalMessage) Encode() []byte { return append(bs, bt...) } -func (c *DisposalMessage) 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/drop_message.go similarity index 60% rename from pkg/message/delete2_message.go rename to pkg/message/drop_message.go index 85066a7b..52dbabe1 100644 --- a/pkg/message/delete2_message.go +++ b/pkg/message/drop_message.go @@ -4,24 +4,26 @@ import ( "encoding/binary" ) -// Requests deletion of all objects under a given prefix. -type DisposalPrefixMessage struct { // Seg port +// 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 = &DisposalPrefixMessage{} +var _ ProtoMessage = &DropMessage{} -func BuildDisposalPrefixMessage(prefix string, confirm bool, garbage bool) *DisposalPrefixMessage { - return &DisposalPrefixMessage{ +func BuildDropMessage(prefix string, confirm bool, garbage bool) *DropMessage { + return &DropMessage{ Prefix: prefix, Confirm: confirm, Garbage: garbage, } } -func (c *DisposalPrefixMessage) Encode() []byte { +func (c *DropMessage) Encode() []byte { bt := []byte{ byte(MessageTypeDelete2), 0, @@ -45,7 +47,7 @@ func (c *DisposalPrefixMessage) Encode() []byte { return append(bs, bt...) } -func (c *DisposalPrefixMessage) Decode(body []byte) { +func (c *DropMessage) Decode(body []byte) { if body[1] == 1 { c.Confirm = true } diff --git a/pkg/message/message_test.go b/pkg/message/message_test.go index 2983c1fb..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.BuildDisposalMessage("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.DisposalMessage{} + 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.BuildDisposalPrefixMessage("trash/prefix", true, false) + msg := message.BuildDropMessage("trash/prefix", true, false) body := msg.Encode() assert.Equal(body[8], byte(message.MessageTypeDelete2)) - msg2 := message.DisposalPrefixMessage{} + 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 153f783d..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 { - HandleDisposalGarbage(message.DisposalMessage) error - HandleDeleteFile(message.DisposalMessage) error + HandleGarbageCleanup(message.CleanupMessage) error + HandleFileDeletion(message.CleanupMessage) error HandleUntrashifyFile(message.UntrashifyMessage) error } @@ -104,7 +104,7 @@ func (dh *BasicGarbageMgr) HandleUntrashifyFile(msg message.UntrashifyMessage) e return nil } -type fileDisposal struct { +type garbageCleanupStrategy struct { workerCount int defaultWorkerCount int failedActionMsg string @@ -115,9 +115,9 @@ type fileDisposal struct { } // Deletes garbage files immediately (hard delete, CrazyDrop). -func (dh *BasicGarbageMgr) deleteDisposal() fileDisposal { +func (dh *BasicGarbageMgr) hardDeleteStrategy() garbageCleanupStrategy { s := dh.StorageInterractor - return fileDisposal{ + return garbageCleanupStrategy{ workerCount: dh.Cnf.TrashDeleteWorkers, defaultWorkerCount: config.DefaultTrashDeleteWorkers, failedActionMsg: "failed to delete some files", @@ -139,9 +139,9 @@ func (dh *BasicGarbageMgr) deleteDisposal() fileDisposal { } // Moves garbage files to the trash prefix (soft delete). -func (dh *BasicGarbageMgr) moveDisposal(segnum int) fileDisposal { +func (dh *BasicGarbageMgr) softDeleteStrategy(segnum int) garbageCleanupStrategy { s := dh.StorageInterractor - return fileDisposal{ + return garbageCleanupStrategy{ workerCount: dh.Cnf.TrashMoveWorkers, defaultWorkerCount: config.DefaultTrashMoveWorkers, failedActionMsg: "failed to move some files", @@ -168,18 +168,18 @@ func (dh *BasicGarbageMgr) moveDisposal(segnum int) fileDisposal { } } -// chooseDisposal picks how garbage files should be disposed of for this request. -func (dh *BasicGarbageMgr) chooseDisposal(msg message.DisposalMessage) fileDisposal { +// chooseDeletionStrategy selects soft or hard deletion for this request. +func (dh *BasicGarbageMgr) chooseDeletionStrategy(msg message.CleanupMessage) garbageCleanupStrategy { if msg.CrazyDrop { - return dh.deleteDisposal() + return dh.hardDeleteStrategy() } - return dh.moveDisposal(int(msg.Segnum)) + return dh.softDeleteStrategy(int(msg.Segnum)) } -func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.DisposalMessage) error { +func (dh *BasicGarbageMgr) CleanupGarbageInBucket(bucket string, msg message.CleanupMessage) error { start := time.Now() t := metrics.NewDeleteOpTracker(bucket, "DELETE_GARBAGE") - disposal := dh.chooseDisposal(msg) + strategy := dh.chooseDeletionStrategy(msg) fileList, err := dh.ListGarbageFiles(bucket, msg) if err != nil { @@ -192,7 +192,7 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di ylogger.Zero.Info().Str("bucket", bucket).Int("files", len(fileList)).Int("uploads", len(uploads)).Msg("garbage delete started") for _, file := range fileList { - disposal.logCandidate(bucket, file) + strategy.logCandidate(bucket, file) } for _, upload := range uploads { ylogger.Zero.Info().Str("bucket", bucket).Str("uploadId", upload).Msg("upload will be aborted") @@ -210,21 +210,21 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di fileList, err = dh.OperateFilesParallel( bucket, batch, - disposal.workerCount, - disposal.defaultWorkerCount, + strategy.workerCount, + strategy.defaultWorkerCount, func(file *object.ObjectInfo) error { - opErr := disposal.execute(bucket, file) + opErr := strategy.execute(bucket, file) if opErr == nil { - disposal.complete(t, file.Size) + strategy.complete(t, file.Size) } return opErr }, - disposal.failedActionMsg, + strategy.failedActionMsg, t, ) deleted += len(batch) - len(fileList) if err != nil { - ylogger.Zero.Error().Str("bucket", bucket).AnErr("err", err).Msg(disposal.failedActionMsg) + ylogger.Zero.Error().Str("bucket", bucket).AnErr("err", err).Msg(strategy.failedActionMsg) } } @@ -232,9 +232,9 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di for _, file := range fileList { t.RecordFailed(file.Size) } - ylogger.Zero.Error().Str("bucket", bucket).Int("failed files count", len(fileList)).Msg(disposal.failedFilesMsg) - ylogger.Zero.Error().Str("bucket", bucket).Any("failed files", fileList).Msg(disposal.failedActionMsg) - return errors.Wrap(err, disposal.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 { @@ -248,16 +248,16 @@ func (dh *BasicGarbageMgr) DisposalGarbageInBucket(bucket string, msg message.Di return nil } -func (dh *BasicGarbageMgr) HandleDisposalGarbage(msg message.DisposalMessage) error { +func (dh *BasicGarbageMgr) HandleGarbageCleanup(msg message.CleanupMessage) error { for _, b := range dh.StorageInterractor.ListBuckets() { - if err := dh.DisposalGarbageInBucket(b, msg); err != nil { + if err := dh.CleanupGarbageInBucket(b, msg); err != nil { return err } } return nil } -func (dh *BasicGarbageMgr) ListDeletePrefixFiles(bucket string, msg message.DisposalPrefixMessage) ([]*object.ObjectInfo, error) { +func (dh *BasicGarbageMgr) ListDeletePrefixFiles(bucket string, msg message.DropMessage) ([]*object.ObjectInfo, error) { // Get first backup lsn var err error @@ -364,7 +364,7 @@ func (dh *BasicGarbageMgr) DeleteGarbageParallel(bucket string, fileList []*obje ) } -func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.DisposalPrefixMessage) error { +func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.DropMessage) error { start := time.Now() t := metrics.NewDeleteOpTracker(bucket, "DELETE_PREFIX") @@ -435,7 +435,7 @@ func (dh *BasicGarbageMgr) DeletePrefixInBucket(bucket string, msg message.Dispo return nil } -func (dh *BasicGarbageMgr) HandleDeletePrefix(msg message.DisposalPrefixMessage) 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 @@ -445,7 +445,7 @@ func (dh *BasicGarbageMgr) HandleDeletePrefix(msg message.DisposalPrefixMessage) } // Delete a single external object by exact name. -func (dh *BasicGarbageMgr) HandleDeleteFile(msg message.DisposalMessage) error { +func (dh *BasicGarbageMgr) HandleFileDeletion(msg message.CleanupMessage) error { if !msg.Confirm { return nil } @@ -459,7 +459,7 @@ func (dh *BasicGarbageMgr) HandleDeleteFile(msg message.DisposalMessage) error { return nil } -func (dh *BasicGarbageMgr) ListGarbageFiles(bucket string, msg message.DisposalMessage) ([]*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 24c4c451..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.DisposalMessage{ + 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.DisposalMessage{ + 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.DisposalMessage{ + 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.DisposalMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -272,7 +272,7 @@ func TestTrashPathConversion(t *testing.T) { func TestListDeletePrefixFiles(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DisposalPrefixMessage{ + msg := message.DropMessage{ Prefix: "trash", Garbage: true, Confirm: true, @@ -304,7 +304,7 @@ func TestListDeletePrefixFiles(t *testing.T) { func TestDeletePrefixInBucketDeletesAllFilesOnceInParallelGarbagePass(t *testing.T) { ctrl := gomock.NewController(t) - msg := message.DisposalPrefixMessage{ + 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.DisposalPrefixMessage{ + 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.DisposalPrefixMessage{ + 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.DisposalPrefixMessage{ + 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.DisposalPrefixMessage{ + 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.DisposalMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -538,7 +538,7 @@ func TestDeleteGarbageInBucketMovesObjectsWhenCrazyDropDisabled(t *testing.T) { Cnf: &config.InstanceConfig().VacuumCnf, } - err := handler.DisposalGarbageInBucket("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.DisposalMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -700,7 +700,7 @@ func TestDeleteGarbageInBucketRecordsMetrics(t *testing.T) { Cnf: &config.InstanceConfig().VacuumCnf, } - err := handler.DisposalGarbageInBucket(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.DisposalPrefixMessage{ + 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.DisposalPrefixMessage{ + 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.DisposalMessage{ + msg := message.CleanupMessage{ Name: "path", Port: 6000, Segnum: 0, @@ -869,7 +869,7 @@ func TestDeleteGarbageInBucketRetriesFailedTrashMoves(t *testing.T) { Cnf: &config.InstanceConfig().VacuumCnf, } - err := handler.DisposalGarbageInBucket("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 57442678..ed1bbccf 100644 --- a/pkg/proc/interaction.go +++ b/pkg/proc/interaction.go @@ -452,8 +452,8 @@ func (*ProtoMgrImpl) ProcessCopyExtended( return nil } -func (*ProtoMgrImpl) ProcessDisposalExtended( - msg message.DisposalMessage, +func (*ProtoMgrImpl) ProcessCleanupExtended( + msg message.CleanupMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, @@ -473,17 +473,17 @@ func (*ProtoMgrImpl) ProcessDisposalExtended( var ( logMsg string successMsg string - handleDelete func(msg message.DisposalMessage) error + handleDelete func(msg message.CleanupMessage) error ) if msg.Garbage { logMsg = "requested to perform external storage VACUUM" successMsg = "Deleted garbage successfully" - handleDelete = dh.HandleDisposalGarbage + handleDelete = dh.HandleGarbageCleanup } else { logMsg = "requested to remove external chunk" successMsg = "Deleted chunk successfully" - handleDelete = dh.HandleDeleteFile + handleDelete = dh.HandleFileDeletion } ylogger.Zero.Debug(). @@ -511,8 +511,8 @@ func (*ProtoMgrImpl) ProcessDisposalExtended( return nil } -func (*ProtoMgrImpl) ProcessDisposalPrefixExtended( - msg message.DisposalPrefixMessage, +func (*ProtoMgrImpl) ProcessDropExtended( + msg message.DropMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, @@ -854,16 +854,16 @@ func ProcConn( case message.MessageTypeDelete: // receive message - msg := message.DisposalMessage{} + msg := message.CleanupMessage{} msg.Decode(body) - err := m.ProcessDisposalExtended(msg, s, bs, ycl, cnf) + err := m.ProcessCleanupExtended(msg, s, bs, ycl, cnf) if err != nil { return err } case message.MessageTypeDelete2: - msg := message.DisposalPrefixMessage{} + msg := message.DropMessage{} msg.Decode(body) - err := m.ProcessDisposalPrefixExtended(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 3221130e..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 */ - ProcessDisposalExtended( - msg message.DisposalMessage, + ProcessCleanupExtended( + msg message.CleanupMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, cnf *config.Vacuum) error - ProcessDisposalPrefixExtended( - msg message.DisposalPrefixMessage, + ProcessDropExtended( + msg message.DropMessage, s storage.StorageInteractor, bs storage.StorageInteractor, ycl client.YproxyClient, From e1e530ef146d78b2e7316294a4e80d45e0688b32 Mon Sep 17 00:00:00 2001 From: Vlasdislav Date: Fri, 4 Sep 2026 21:43:12 +0300 Subject: [PATCH 4/4] Refactor: make fmt --- cmd/client/main.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cmd/client/main.go b/cmd/client/main.go index be91dd0a..249e6e93 100644 --- a/cmd/client/main.go +++ b/cmd/client/main.go @@ -404,7 +404,7 @@ var copyCmd = &cobra.Command{ var deleteCmd = &cobra.Command{ Use: "delete", Short: "remove an object or collect garbage", - RunE: Runner(sendCleanupRequest), + RunE: Runner(sendCleanupRequest), Args: cobra.ExactArgs(1), } @@ -437,7 +437,7 @@ var goolCmd = &cobra.Command{ var deletePrefixCmd = &cobra.Command{ Use: "deletePrefix", Short: "physically delete objects by prefix", - RunE: Runner(sendDropRequest), + RunE: Runner(sendDropRequest), Args: cobra.ExactArgs(1), // name_prefix }