| 594 | } |
| 595 | |
| 596 | func (c *BlocksCleaner) cleanUser(ctx context.Context, userLogger log.Logger, userBucket objstore.InstrumentedBucket, userID string, firstRun bool) (returnErr error) { |
| 597 | c.blocksMarkedForDeletion.WithLabelValues(userID, reasonValueRetention) |
| 598 | startTime := time.Now() |
| 599 | |
| 600 | bucketIndexDeleted := false |
| 601 | |
| 602 | level.Info(userLogger).Log("msg", "started blocks cleanup and maintenance") |
| 603 | defer func() { |
| 604 | if returnErr != nil { |
| 605 | level.Warn(userLogger).Log("msg", "failed blocks cleanup and maintenance", "err", returnErr) |
| 606 | } else { |
| 607 | level.Info(userLogger).Log("msg", "completed blocks cleanup and maintenance", "duration", time.Since(startTime), "duration_ms", time.Since(startTime).Milliseconds()) |
| 608 | } |
| 609 | c.tenantCleanDuration.WithLabelValues(userID).Set(time.Since(startTime).Seconds()) |
| 610 | }() |
| 611 | |
| 612 | if c.cfg.ShardingStrategy == util.ShardingStrategyShuffle && c.cfg.CompactionStrategy == util.CompactionStrategyPartitioning { |
| 613 | begin := time.Now() |
| 614 | c.cleanPartitionedGroupInfo(ctx, userBucket, userLogger, userID) |
| 615 | level.Info(userLogger).Log("msg", "finish cleaning partitioned group info files", "duration", time.Since(begin), "duration_ms", time.Since(begin).Milliseconds()) |
| 616 | } |
| 617 | |
| 618 | // Migrate block deletion marks to the global markers location. This operation is a best-effort. |
| 619 | if firstRun && c.cfg.BlockDeletionMarksMigrationEnabled { |
| 620 | if err := bucketindex.MigrateBlockDeletionMarksToGlobalLocation(ctx, c.bucketClient, userID, c.cfgProvider); err != nil { |
| 621 | level.Warn(userLogger).Log("msg", "failed to migrate block deletion marks to the global markers location", "err", err) |
| 622 | } else { |
| 623 | level.Info(userLogger).Log("msg", "migrated block deletion marks to the global markers location") |
| 624 | } |
| 625 | } |
| 626 | |
| 627 | // Reading bucket index sync stats |
| 628 | idxs, err := bucketindex.ReadSyncStatus(ctx, c.bucketClient, userID, userLogger) |
| 629 | |
| 630 | if err != nil { |
| 631 | level.Warn(userLogger).Log("msg", "error reading the bucket index status", "err", err) |
| 632 | idxs = bucketindex.Status{Version: bucketindex.SyncStatusFileVersion, NonQueryableReason: bucketindex.Unknown} |
| 633 | } |
| 634 | |
| 635 | idxs.Status = bucketindex.Ok |
| 636 | idxs.SyncTime = time.Now().Unix() |
| 637 | |
| 638 | // Read the bucket index. |
| 639 | begin := time.Now() |
| 640 | idx, err := bucketindex.ReadIndex(ctx, c.bucketClient, userID, c.cfgProvider, c.logger) |
| 641 | |
| 642 | defer func() { |
| 643 | if bucketIndexDeleted { |
| 644 | level.Info(userLogger).Log("msg", "deleting bucket index sync status since bucket index is empty") |
| 645 | if err := bucketindex.DeleteIndexSyncStatus(ctx, c.bucketClient, userID); err != nil { |
| 646 | level.Warn(userLogger).Log("msg", "error deleting index sync status when index is empty", "err", err) |
| 647 | } |
| 648 | if err := c.deleteNonDataFiles(ctx, userLogger, userBucket); err != nil { |
| 649 | level.Warn(userLogger).Log("msg", "error deleting non-data files", "err", err) |
| 650 | } |
| 651 | } else { |
| 652 | bucketindex.WriteSyncStatus(ctx, c.bucketClient, userID, idxs, userLogger) |
| 653 | } |