Files
files/internal/jobs/integrity.go

152 lines
3.8 KiB
Go
Raw Permalink Normal View History

2026-09-09 16:42:21 +08:00
package jobs
import (
"context"
"errors"
"fmt"
"sync"
"time"
"git.apinb.com/bsm-sdk/core/logger"
"git.apinb.com/ops/files/internal/config"
"git.apinb.com/ops/files/internal/impl"
"git.apinb.com/ops/files/internal/models"
"git.apinb.com/ops/files/internal/storage"
)
type IntegrityStatus struct {
Running bool `json:"running"`
LastStarted time.Time `json:"last_started,omitempty"`
LastSuccess time.Time `json:"last_success,omitempty"`
LastError string `json:"last_error,omitempty"`
Checked int64 `json:"checked"`
Missing int64 `json:"missing"`
SizeMismatch int64 `json:"size_mismatch"`
ETagMismatch int64 `json:"etag_mismatch"`
}
var integrityRuntime struct {
sync.RWMutex
IntegrityStatus
wg sync.WaitGroup
}
func StartIntegrity(ctx context.Context) error {
integrityRuntime.Lock()
if integrityRuntime.Running {
integrityRuntime.Unlock()
return fmt.Errorf("文件完整性任务已启动")
}
integrityRuntime.Unlock()
if err := runIntegrity(ctx); err != nil {
return fmt.Errorf("首次核对文件完整性失败: %w", err)
}
integrityRuntime.Lock()
integrityRuntime.Running = true
integrityRuntime.Unlock()
integrityRuntime.wg.Add(1)
go func() {
defer integrityRuntime.wg.Done()
defer func() {
integrityRuntime.Lock()
integrityRuntime.Running = false
integrityRuntime.Unlock()
}()
ticker := time.NewTicker(time.Duration(config.Spec.Cleanup.IntegrityIntervalSeconds) * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := runIntegrity(ctx); err != nil && !errors.Is(err, context.Canceled) {
logger.Errorf("stage=files_integrity error=%v", err)
}
}
}
}()
return nil
}
func IntegrityStatusSnapshot() IntegrityStatus {
integrityRuntime.RLock()
defer integrityRuntime.RUnlock()
return integrityRuntime.IntegrityStatus
}
func waitIntegrity(ctx context.Context) error {
done := make(chan struct{})
go func() {
integrityRuntime.wg.Wait()
close(done)
}()
select {
case <-done:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func runIntegrity(ctx context.Context) error {
startedAt := time.Now().UTC()
status := IntegrityStatus{Running: true, LastStarted: startedAt}
integrityRuntime.Lock()
status.LastSuccess = integrityRuntime.LastSuccess
integrityRuntime.IntegrityStatus = status
integrityRuntime.Unlock()
var lastID uint
for {
var fileObjects []models.FileObject
result := impl.DBService.WithContext(ctx).
Where("status = ? AND id > ?", models.FileStatusReady, lastID).
Order("id ASC").Limit(config.Spec.Cleanup.IntegrityBatchSize).Find(&fileObjects)
if result.Error != nil {
return finishIntegrity(status, result.Error)
}
if len(fileObjects) == 0 {
break
}
for _, fileObject := range fileObjects {
if err := ctx.Err(); err != nil {
return finishIntegrity(status, err)
}
status.Checked++
objectInfo, err := impl.StorageService.Stat(ctx, fileObject.ObjectKey)
if err != nil {
if errors.Is(err, storage.ErrObjectNotFound) {
status.Missing++
continue
}
return finishIntegrity(status, err)
}
if objectInfo.Size != fileObject.ActualSize {
status.SizeMismatch++
}
if objectInfo.StorageETag != fileObject.StorageETag {
status.ETagMismatch++
}
}
lastID = fileObjects[len(fileObjects)-1].ID
if len(fileObjects) < config.Spec.Cleanup.IntegrityBatchSize {
break
}
}
status.LastSuccess = time.Now().UTC()
status.Running = integrityRuntime.Running
integrityRuntime.Lock()
integrityRuntime.IntegrityStatus = status
integrityRuntime.Unlock()
return nil
}
func finishIntegrity(status IntegrityStatus, err error) error {
status.LastError = err.Error()
status.Running = integrityRuntime.Running
integrityRuntime.Lock()
integrityRuntime.IntegrityStatus = status
integrityRuntime.Unlock()
return err
}