Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 10 additions & 10 deletions internal/removequeue.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,11 @@ type (
}

RemoveOperation struct {
Status Status `json:"status"`
Removed *collections.Set[string] `json:"removed"` // Object names
Failed *collections.Set[string] `json:"failed"` // Object names
Errors map[string]error `json:"errors"`
LastUpdated time.Time `json:"-"`
Status Status `json:"status"`
Removed *collections.Set[string] `json:"removed"` // Object names
Failed *collections.Set[string] `json:"failed"` // Object names
Errors map[string]string `json:"errors"`
LastUpdated time.Time `json:"-"`
}

Status string
Expand Down Expand Up @@ -57,7 +57,7 @@ func (q *RemoveQueue) StartOperation(guildId uint64) error {
Status: StatusInProgress,
Removed: collections.NewSet[string](),
Failed: collections.NewSet[string](),
Errors: make(map[string]error),
Errors: make(map[string]string),
LastUpdated: time.Now(),
}

Expand Down Expand Up @@ -113,7 +113,7 @@ func (q *RemoveQueue) AddError(guildId uint64, objectName string, err error) {
defer q.mu.Unlock()

if operation, ok := q.queue[guildId]; ok {
operation.Errors[objectName] = err
operation.Errors[objectName] = err.Error()
operation.Failed.Add(objectName)
operation.Removed.Remove(objectName)
operation.LastUpdated = time.Now()
Expand Down Expand Up @@ -168,16 +168,16 @@ func (o RemoveOperation) clone() RemoveOperation {
failed.Add(el)
}

errors := make(map[string]error)
errs := make(map[string]string, len(o.Errors))
for k, v := range o.Errors {
errors[k] = v
errs[k] = v
}

return RemoveOperation{
Status: o.Status,
Removed: removed,
Failed: failed,
Errors: errors,
Errors: errs,
LastUpdated: o.LastUpdated,
}
}
41 changes: 34 additions & 7 deletions pkg/http/purgeguild.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ func (s *Server) purgeGuildHandler(ctx *gin.Context) {
return
}); err != nil {
s.Logger.Error("Failed to fetch objects from database", zap.Error(err), zap.Uint64("guild", guildId))
s.RemoveQueue.AddError(guildId, "database:list", err)
firstErr = err
} else {
for _, obj := range dbObjects {
Expand All @@ -71,6 +72,7 @@ func (s *Server) purgeGuildHandler(ctx *gin.Context) {
bucketId, err := uuid.Parse(bucketIdStr)
if err != nil {
s.Logger.Error("Failed to parse bucket ID", zap.Error(err), zap.String("bucket_id", bucketIdStr), zap.Uint64("guild", guildId))
s.RemoveQueue.AddError(guildId, "bucket:"+bucketIdStr, err)
if firstErr == nil {
firstErr = err
}
Expand All @@ -80,6 +82,7 @@ func (s *Server) purgeGuildHandler(ctx *gin.Context) {
client, err := s.s3Clients.Get(bucketId)
if err != nil {
s.Logger.Error("Failed to get S3 client", zap.Error(err), zap.String("bucket_id", bucketIdStr), zap.Uint64("guild", guildId))
s.RemoveQueue.AddError(guildId, "bucket:"+bucketIdStr, err)
if firstErr == nil {
firstErr = err
}
Expand All @@ -98,16 +101,24 @@ func (s *Server) purgeGuildHandler(ctx *gin.Context) {
objectsCh := make(chan minio.ObjectInfo)

bucketObjectCount := 0
listErrCh := make(chan error, 1)
go func() {
defer close(objectsCh)

var listErr error
for obj := range objectCh {
if obj.Err != nil {
s.Logger.Warn("Error listing object (non-fatal)", zap.Error(obj.Err), zap.Uint64("guild", guildId))
s.Logger.Error("Failed to list object", zap.Error(obj.Err), zap.String("bucket_id", bucketIdStr), zap.Uint64("guild", guildId))
if listErr == nil {
listErr = obj.Err
}
continue
}
bucketObjectCount++
objectsCh <- obj
}

listErrCh <- listErr
s.Logger.Debug("Finished listing objects", zap.Int("count", bucketObjectCount), zap.Uint64("guild", guildId))
}()

Expand All @@ -116,6 +127,7 @@ func (s *Server) purgeGuildHandler(ctx *gin.Context) {
for result := range client.Minio().RemoveObjects(context.Background(), client.BucketName(), objectsCh, minio.RemoveObjectsOptions{}) {
if result.Err != nil {
s.Logger.Error("Failed to remove object", zap.Error(result.Err), zap.String("object", result.ObjectName), zap.Uint64("guild", guildId))
s.RemoveQueue.AddError(guildId, result.ObjectName, result.Err)
if firstErr == nil {
firstErr = result.Err
}
Expand All @@ -125,17 +137,32 @@ func (s *Server) purgeGuildHandler(ctx *gin.Context) {
}
}

// Reading this also joins the lister.
if listErr := <-listErrCh; listErr != nil {
s.RemoveQueue.AddError(guildId, "bucket:"+bucketIdStr, listErr)
if firstErr == nil {
firstErr = listErr
}
}

s.RemoveQueue.AddRemovedObject(guildId, fmt.Sprintf("%s (%d objects)", bucketIdStr, bucketObjectCount))

s.Logger.Debug("Finished deleting from bucket", zap.Int("deleted", bucketDeletedCount), zap.String("bucket_id", bucketIdStr), zap.Uint64("guild", guildId))
}

// Delete database records
if err := s.store.Tx(context.Background(), func(r repository.Repositories) error {
return r.Objects().DeleteByGuild(context.Background(), guildId)
}); err != nil {
s.Logger.Error("Failed to delete objects from database", zap.Error(err), zap.Uint64("guild", guildId))
if firstErr == nil {
if firstErr == nil {
if err := s.store.Tx(context.Background(), func(r repository.Repositories) error {
return r.Objects().DeleteByGuild(context.Background(), guildId)
}); err != nil {
s.Logger.Error("Failed to delete objects from database", zap.Error(err), zap.Uint64("guild", guildId))
s.RemoveQueue.AddError(guildId, "database:delete", err)
firstErr = err
}
} else {
s.Logger.Warn(
"Keeping object records because the S3 purge failed; deleting them would strand the transcripts",
zap.Uint64("guild", guildId),
)
}

if firstErr != nil {
Expand Down
Loading