From f9f82f455aefa38b7a466a6c7a6cfdaaa3c77942 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Thu, 24 Sep 2026 06:52:21 +0300 Subject: [PATCH] fix: avoid skipping orphan file records MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../file_record/file_record_service.go | 15 +- .../file_record/file_record_service_test.go | 138 ++++++++++++++++++ 2 files changed, 150 insertions(+), 3 deletions(-) create mode 100644 internal/service/file_record/file_record_service_test.go diff --git a/internal/service/file_record/file_record_service.go b/internal/service/file_record/file_record_service.go index aa526f014..1ff1152fd 100644 --- a/internal/service/file_record/file_record_service.go +++ b/internal/service/file_record/file_record_service.go @@ -92,10 +92,10 @@ func (fs *FileRecordService) AddFileRecord(ctx context.Context, userID, filePath // CleanOrphanUploadFiles clean orphan upload files func (fs *FileRecordService) CleanOrphanUploadFiles(ctx context.Context) { - page, pageSize := 1, 1000 + pageSize := 1000 for { - fileRecordList, total, err := fs.fileRecordRepo.GetFileRecordPage(ctx, page, pageSize, &entity.FileRecord{ + fileRecordList, total, err := fs.fileRecordRepo.GetFileRecordPage(ctx, 1, pageSize, &entity.FileRecord{ Status: entity.FileRecordStatusAvailable, }) if err != nil { @@ -105,6 +105,7 @@ func (fs *FileRecordService) CleanOrphanUploadFiles(ctx context.Context) { if len(fileRecordList) == 0 || total == 0 { break } + changed := false for _, fileRecord := range fileRecordList { // If this file record created in 48 hours, no need to check if fileRecord.CreatedAt.AddDate(0, 0, 2).After(time.Now()) { @@ -122,6 +123,8 @@ func (fs *FileRecordService) CleanOrphanUploadFiles(ctx context.Context) { } if err := fs.DeleteAndMoveFileRecord(ctx, fileRecord); err != nil { log.Error(err) + } else { + changed = true } continue } @@ -145,6 +148,8 @@ func (fs *FileRecordService) CleanOrphanUploadFiles(ctx context.Context) { fileRecord.ObjectID = lastRevision.ObjectID if err := fs.fileRecordRepo.UpdateFileRecord(ctx, fileRecord); err != nil { log.Errorf("update file record object id error: %v", err) + } else { + changed = true } continue } @@ -152,9 +157,13 @@ func (fs *FileRecordService) CleanOrphanUploadFiles(ctx context.Context) { // Delete and move the file record if err := fs.DeleteAndMoveFileRecord(ctx, fileRecord); err != nil { log.Error(err) + } else { + changed = true } } - page++ + if !changed { + break + } } } diff --git a/internal/service/file_record/file_record_service_test.go b/internal/service/file_record/file_record_service_test.go new file mode 100644 index 000000000..afaf53c29 --- /dev/null +++ b/internal/service/file_record/file_record_service_test.go @@ -0,0 +1,138 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package file_record + +import ( + "context" + "os" + "path/filepath" + "strconv" + "testing" + "time" + + "github.com/apache/answer/internal/base/constant" + "github.com/apache/answer/internal/entity" + "github.com/apache/answer/internal/service/revision" + "github.com/apache/answer/internal/service/service_config" +) + +type testFileRecordRepo struct { + records []*entity.FileRecord +} + +func (r *testFileRecordRepo) AddFileRecord(context.Context, *entity.FileRecord) error { + return nil +} + +func (r *testFileRecordRepo) UpdateFileRecord(_ context.Context, fileRecord *entity.FileRecord) error { + for _, record := range r.records { + if record.ID == fileRecord.ID { + record.ObjectID = fileRecord.ObjectID + } + } + return nil +} + +func (r *testFileRecordRepo) GetFileRecordPage(_ context.Context, page, pageSize int, _ *entity.FileRecord) ([]*entity.FileRecord, int64, error) { + available := make([]*entity.FileRecord, 0, len(r.records)) + for _, record := range r.records { + if record.Status == entity.FileRecordStatusAvailable { + available = append(available, record) + } + } + start := (page - 1) * pageSize + if start >= len(available) { + return []*entity.FileRecord{}, int64(len(available)), nil + } + end := start + pageSize + if end > len(available) { + end = len(available) + } + return available[start:end], int64(len(available)), nil +} + +func (r *testFileRecordRepo) DeleteFileRecord(_ context.Context, id int) error { + for _, record := range r.records { + if record.ID == id { + record.Status = entity.FileRecordStatusDeleted + } + } + return nil +} + +func (r *testFileRecordRepo) GetFileRecordByURL(context.Context, string) (*entity.FileRecord, error) { + return nil, nil +} + +type testRevisionRepo struct { + revision.RevisionRepo +} + +func (testRevisionRepo) GetLastRevisionByObjectID(context.Context, string) (*entity.Revision, bool, error) { + return nil, false, nil +} + +func (testRevisionRepo) GetLastRevisionByFileURL(context.Context, string) (*entity.Revision, bool, error) { + return nil, false, nil +} + +func TestCleanOrphanUploadFilesDoesNotSkipRecordsWhenStatusChanges(t *testing.T) { + const recordCount = 1001 + uploadPath := t.TempDir() + if err := os.MkdirAll(filepath.Join(uploadPath, "tmp"), 0o755); err != nil { + t.Fatal(err) + } + if err := os.MkdirAll(filepath.Join(uploadPath, constant.DeletedSubPath), 0o755); err != nil { + t.Fatal(err) + } + + records := make([]*entity.FileRecord, 0, recordCount) + createdAt := time.Now().AddDate(0, 0, -3) + for id := 1; id <= recordCount; id++ { + filePath := filepath.Join("tmp", "orphan-"+strconv.Itoa(id)) + if err := os.WriteFile(filepath.Join(uploadPath, filePath), []byte("orphan"), 0o600); err != nil { + t.Fatal(err) + } + records = append(records, &entity.FileRecord{ + ID: id, + CreatedAt: createdAt, + FilePath: filePath, + FileURL: "/" + filePath, + ObjectID: "1", + Status: entity.FileRecordStatusAvailable, + }) + } + + repo := &testFileRecordRepo{records: records} + service := NewFileRecordService( + repo, + testRevisionRepo{}, + &service_config.ServiceConfig{UploadPath: uploadPath}, + nil, + nil, + ) + service.CleanOrphanUploadFiles(context.Background()) + + for _, record := range records { + if record.Status != entity.FileRecordStatusDeleted { + t.Fatalf("orphan file record %d was not deleted", record.ID) + } + } +}