package documentclean import ( "context" "database/sql" "errors" "fmt" "log/slog" "queryorchestration/internal/database/repository" documenttypes "queryorchestration/internal/document/types" "queryorchestration/internal/serviceconfig/build" "queryorchestration/internal/serviceconfig/objectstore" "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/google/uuid" ) type CleanParams struct { ID uuid.UUID Hash string Bucket string Key objectstore.BucketKey } type ExecuteCleanResponse struct { key *objectstore.BucketKey bucket *string mimetype *documenttypes.MimeType hash *string failReason *documenttypes.InvalidDocumentReason } func (s *Service) executeCleanTasks(ctx context.Context, params *CleanParams) (*ExecuteCleanResponse, error) { keyStr := params.Key.String() out, err := s.cfg.GetStoreClient().HeadObject(ctx, &s3.HeadObjectInput{ Bucket: ¶ms.Bucket, Key: &keyStr, IfMatch: ¶ms.Hash, }) if err != nil { return nil, err } length := int64(1) if out.ContentLength != nil { length = *out.ContentLength } mimeType, err := s.getAcceptedMimeType(ctx, params, out.ContentType, length) if err != nil { return nil, err } else if mimeType == documenttypes.MimeTypeInvalid { reason := documenttypes.InvalidDocumentMimeType return &ExecuteCleanResponse{ failReason: &reason, }, nil } content, err := s.getContent(ctx, params, mimeType) if err != nil { return nil, err } corruptReason := content.IsCorrupt(ctx) if corruptReason != nil { return &ExecuteCleanResponse{ failReason: corruptReason, }, nil } return &ExecuteCleanResponse{ mimetype: &mimeType, bucket: ¶ms.Bucket, key: ¶ms.Key, hash: ¶ms.Hash, }, nil } func (s *Service) clean(ctx context.Context, id uuid.UUID) error { slog.Debug("cleaning document", "id", id.String()) docId := id doc, err := s.cfg.GetDBQueries().GetDocumentSummary(ctx, docId) if err != nil { return err } entry, err := s.cfg.GetDBQueries().GetDocumentEntry(ctx, docId) if err != nil { return err } key, err := objectstore.ParseBucketKey(entry.Key) if err != nil { return err } out, err := s.executeCleanTasks(ctx, &CleanParams{ ID: id, Hash: doc.Hash, Bucket: entry.Bucket, Key: key, }) if err != nil { return err } err = s.storeClean(ctx, id, out) if err != nil { return err } return nil } func (s *Service) storeClean(ctx context.Context, id uuid.UUID, out *ExecuteCleanResponse) error { docId := id version := build.GetVersionUnixTimestamp() params := &repository.AddDocumentCleanParams{ Documentid: docId, } if out.failReason != nil { slog.Info("Failed document", "id", id, "reason", *out.failReason) params.Fail = documenttypes.ToDBNullFailType(*out.failReason) } else { params.Bucket = out.bucket key := out.key.String() params.Key = &key params.Hash = out.hash params.Mimetype = documenttypes.ToDBNullMimeType(*out.mimetype) } err := s.cfg.ExecuteDBTransaction(ctx, func(ctx context.Context, q *repository.Queries) error { lastClean, err := q.GetMostRecentDocumentCleanEntry(ctx, docId) if err != nil && !errors.Is(err, sql.ErrNoRows) { return err } newID := lastClean.ID if s.isNewClean(lastClean, out) { id, err := q.AddDocumentClean(ctx, params) if err != nil { return err } newID = id } err = q.AddDocumentCleanEntry(ctx, &repository.AddDocumentCleanEntryParams{ Cleanid: newID, Version: version, }) if err != nil { return err } return nil }) if err != nil { return err } if out.failReason != nil { return fmt.Errorf("%s", *out.failReason) } return nil } func (s *Service) isNewClean(lastEntry *repository.GetMostRecentDocumentCleanEntryRow, out *ExecuteCleanResponse) bool { if lastEntry == nil { return true } if out.failReason != nil { return *out.failReason != documenttypes.ParseDBNullFailType(lastEntry.Fail) } if lastEntry.Fail.Valid { return true } if out.bucket == nil || out.key == nil || out.mimetype == nil || lastEntry.Bucket == nil || lastEntry.Key == nil { return true } return !(*out.bucket == *lastEntry.Bucket && out.key.String() == *lastEntry.Key && *out.mimetype == documenttypes.ParseDBNullMimeType(lastEntry.Mimetype)) }