| 378 | } |
| 379 | |
| 380 | func (c *BlocksCleaner) cleanDeletedUsers(ctx context.Context, users []string) error { |
| 381 | level.Info(c.logger).Log("msg", "started blocks cleanup and maintenance for deleted users") |
| 382 | c.runsStarted.WithLabelValues(deletedStatus).Inc() |
| 383 | |
| 384 | return concurrency.ForEachUser(ctx, users, c.cfg.CleanupConcurrency, func(ctx context.Context, userID string) error { |
| 385 | userLogger := util_log.WithUserID(userID, c.logger) |
| 386 | userBucket := bucket.NewUserBucketClient(userID, c.bucketClient, c.cfgProvider) |
| 387 | visitMarkerManager, isVisited, err := c.obtainVisitMarkerManager(ctx, userLogger, userBucket) |
| 388 | if err != nil { |
| 389 | return err |
| 390 | } |
| 391 | if isVisited { |
| 392 | return nil |
| 393 | } |
| 394 | errChan := make(chan error, 1) |
| 395 | go visitMarkerManager.HeartBeat(ctx, errChan, c.cleanerVisitMarkerFileUpdateInterval, true) |
| 396 | defer func() { |
| 397 | errChan <- nil |
| 398 | }() |
| 399 | return errors.Wrapf(c.deleteUserMarkedForDeletion(ctx, userLogger, userBucket, userID), "failed to delete user marked for deletion: %s", userID) |
| 400 | }) |
| 401 | } |
| 402 | |
| 403 | func (c *BlocksCleaner) scanUsers(ctx context.Context) ([]string, []string, error) { |
| 404 | active, deleting, deleted, err := c.usersScanner.ScanUsers(ctx) |