Skip to content

Commit c92b22f

Browse files
committed
compact: avoid cleaner deadlock on cancellation
Signed-off-by: yuchen-db <yw3642@gmail.com>
1 parent 51afa7e commit c92b22f

2 files changed

Lines changed: 58 additions & 5 deletions

File tree

pkg/compact/blocks_cleaner.go

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,8 @@ import (
2020
"github.com/thanos-io/thanos/pkg/errutil"
2121
)
2222

23+
const deleteMarkedBlocksConcurrency = 32
24+
2325
// BlocksCleaner is a struct that deletes blocks from bucket which are marked for deletion.
2426
type BlocksCleaner struct {
2527
logger log.Logger
@@ -45,8 +47,6 @@ func NewBlocksCleaner(logger log.Logger, bkt objstore.Bucket, ignoreDeletionMark
4547
// DeleteMarkedBlocks uses ignoreDeletionMarkFilter to gather the blocks that are marked for deletion and deletes those
4648
// if older than given deleteDelay.
4749
func (s *BlocksCleaner) DeleteMarkedBlocks(ctx context.Context) (map[ulid.ULID]struct{}, error) {
48-
const conc = 32
49-
5050
level.Info(s.logger).Log("msg", "started cleaning of blocks marked for deletion")
5151

5252
var (
@@ -55,10 +55,10 @@ func (s *BlocksCleaner) DeleteMarkedBlocks(ctx context.Context) (map[ulid.ULID]s
5555
deletedBlocks = make(map[ulid.ULID]struct{}, 0)
5656
deletionMarkMap = s.ignoreDeletionMarkFilter.DeletionMarkBlocks()
5757
wg sync.WaitGroup
58-
dm = make(chan *metadata.DeletionMark, conc)
58+
dm = make(chan *metadata.DeletionMark, deleteMarkedBlocksConcurrency)
5959
)
6060

61-
for range conc {
61+
for range deleteMarkedBlocksConcurrency {
6262
wg.Go(func() {
6363
for deletionMark := range dm {
6464
if ctx.Err() != nil {
@@ -82,8 +82,13 @@ func (s *BlocksCleaner) DeleteMarkedBlocks(ctx context.Context) (map[ulid.ULID]s
8282
})
8383
}
8484

85+
enqueue:
8586
for _, deletionMark := range deletionMarkMap {
86-
dm <- deletionMark
87+
select {
88+
case dm <- deletionMark:
89+
case <-ctx.Done():
90+
break enqueue
91+
}
8792
}
8893
close(dm)
8994
wg.Wait()

pkg/compact/compact_test.go

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -749,6 +749,54 @@ func TestGarbageCollect_FilterRace(t *testing.T) {
749749
}
750750
}
751751

752+
func TestBlocksCleanerCancellationWhileEnqueuing(t *testing.T) {
753+
bkt := objstore.WithNoopInstr(objstore.NewInMemBucket())
754+
deletionMarkFilter := block.NewIgnoreDeletionMarkFilter(log.NewNopLogger(), bkt, 0, 1)
755+
deletionMarkCount := 2*deleteMarkedBlocksConcurrency + 1
756+
metas := make(map[ulid.ULID]*metadata.Meta, deletionMarkCount)
757+
758+
for i := range deletionMarkCount {
759+
id := ulid.MustNew(uint64(i+1), bytes.NewReader(make([]byte, 10)))
760+
metas[id] = &metadata.Meta{}
761+
762+
var buf bytes.Buffer
763+
testutil.Ok(t, json.NewEncoder(&buf).Encode(&metadata.DeletionMark{
764+
ID: id,
765+
DeletionTime: 0,
766+
Version: 1,
767+
}))
768+
testutil.Ok(t, bkt.Upload(context.Background(), path.Join(id.String(), metadata.DeletionMarkFilename), &buf))
769+
}
770+
771+
synced := promauto.With(nil).NewGaugeVec(prometheus.GaugeOpts{Name: "test_blocks_cleaner_synced"}, []string{"state"})
772+
testutil.Ok(t, deletionMarkFilter.Filter(context.Background(), metas, synced, nil))
773+
774+
cleaner := NewBlocksCleaner(
775+
log.NewNopLogger(),
776+
bkt,
777+
deletionMarkFilter,
778+
0,
779+
promauto.With(nil).NewCounter(prometheus.CounterOpts{Name: "test_blocks_cleaner_cleaned"}),
780+
promauto.With(nil).NewCounter(prometheus.CounterOpts{Name: "test_blocks_cleaner_failures"}),
781+
)
782+
783+
ctx, cancel := context.WithCancel(context.Background())
784+
cancel()
785+
done := make(chan error, 1)
786+
go func() {
787+
_, err := cleaner.DeleteMarkedBlocks(ctx)
788+
done <- err
789+
}()
790+
791+
select {
792+
case err := <-done:
793+
testutil.NotOk(t, err)
794+
testutil.Assert(t, errors.Is(err, context.Canceled), "expected context cancellation, got %v", err)
795+
case <-time.After(time.Second):
796+
t.Fatal("DeleteMarkedBlocks blocked after cancellation")
797+
}
798+
}
799+
752800
func TestCompactExtensions(t *testing.T) {
753801
t.Parallel()
754802

0 commit comments

Comments
 (0)