From 09031da204c4dc5d3e8a0ccec87414f7f57c739f Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Sat, 19 Sep 2026 01:32:48 +0200 Subject: [PATCH] Add storage-policy blob audit workflow for orphan objects Physical blobs can linger after failed transfers or inconsistent deletes with no DB entity referencing them. Adds an admin-triggered blob_audit task that lists every object under the policy's blob prefix, diffs them against entity sources for that policy, and optionally deletes the orphans in bounded batches. Orphan lists are capped at 5000 entries in task state with a truncation flag; failed deletions are counted and left for the next audit run. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- frontend/public/locales/en-US/dashboard.json | 7 + frontend/public/locales/zh-CN/dashboard.json | 7 + frontend/src/api/api.ts | 18 ++ frontend/src/api/workflow.ts | 6 + .../Admin/StoragePolicy/BlobAuditDialog.tsx | 84 +++++ .../Admin/StoragePolicy/StoragePolicyCard.tsx | 24 +- .../src/component/Admin/Task/TaskContent.tsx | 5 + pkg/filemanager/workflows/blob_audit.go | 301 ++++++++++++++++++ pkg/filemanager/workflows/blob_audit_test.go | 23 ++ pkg/queue/task.go | 1 + routers/controllers/file.go | 15 + routers/router.go | 7 + service/explorer/workflows.go | 34 ++ 13 files changed, 528 insertions(+), 4 deletions(-) create mode 100644 frontend/src/component/Admin/StoragePolicy/BlobAuditDialog.tsx create mode 100644 pkg/filemanager/workflows/blob_audit.go create mode 100644 pkg/filemanager/workflows/blob_audit_test.go diff --git a/frontend/public/locales/en-US/dashboard.json b/frontend/public/locales/en-US/dashboard.json index 653352c2..43356922 100644 --- a/frontend/public/locales/en-US/dashboard.json +++ b/frontend/public/locales/en-US/dashboard.json @@ -988,6 +988,12 @@ "addXStoragePolicy": "Add {{type}} storage policy", "loadSummary": "Load summary", "policySummary": "{{count}} file Blobs ({{size}})", + "blobAudit": "Audit blobs: {{name}}", + "blobAuditDes": "Scans the storage backend and lists physical objects not referenced by any file record. Useful for finding ghost blobs left by failed transfers.", + "blobAuditDelete": "Delete orphaned blobs", + "blobAuditDeleteWarning": "Orphaned objects will be permanently removed from the storage backend. Files still being uploaded may appear as orphans — only enable this when you are sure no uploads are in progress.", + "blobAuditStart": "Start audit", + "blobAuditSubmitted": "Blob audit task queued.", "sharp": "#", "name": "Name", "type": "Type", @@ -1588,6 +1594,7 @@ "fullTextCopy": "Copy full-text index for file <0>#{{fileID}}", "fullTextChangeOwner": "Change index owner for file <0>#{{fileID}}", "fullTextDelete": "Delete full-text index for {{count}} files", + "blobAudit": "Audit blobs on storage policy #{{policyID}} ({{count}} orphans)", "type": "Type", "node": "Distributed node", "createdBy": "Created by", diff --git a/frontend/public/locales/zh-CN/dashboard.json b/frontend/public/locales/zh-CN/dashboard.json index a80a153f..944d5352 100644 --- a/frontend/public/locales/zh-CN/dashboard.json +++ b/frontend/public/locales/zh-CN/dashboard.json @@ -988,6 +988,12 @@ "addXStoragePolicy": "添加 {{type}} 存储策略", "loadSummary": "加载统计数据", "policySummary": "{{count}} 个文件 Blob ({{size}})", + "blobAudit": "Blob 审计:{{name}}", + "blobAuditDes": "扫描存储后端,列出没有任何文件记录引用的物理对象,用于清理传输失败遗留的幽灵 Blob。", + "blobAuditDelete": "删除孤儿 Blob", + "blobAuditDeleteWarning": "孤儿对象将被从存储后端永久删除。正在上传中的文件可能会被误判为孤儿——请仅在确认没有上传任务进行时启用。", + "blobAuditStart": "开始审计", + "blobAuditSubmitted": "Blob 审计任务已加入队列。", "sharp": "#", "name": "名称", "type": "类型", @@ -1587,6 +1593,7 @@ "fullTextCopy": "复制文件 <0>#{{fileID}} 的全文索引", "fullTextChangeOwner": "变更文件 <0>#{{fileID}} 的索引所有者", "fullTextDelete": "删除 {{count}} 个文件的全文索引", + "blobAudit": "审计存储策略 #{{policyID}} 的 Blob({{count}} 个孤儿对象)", "type": "类型", "node": "处理节点", "createdBy": "创建者", diff --git a/frontend/src/api/api.ts b/frontend/src/api/api.ts index e8900faf..f992197b 100644 --- a/frontend/src/api/api.ts +++ b/frontend/src/api/api.ts @@ -106,6 +106,7 @@ import { DownloadWorkflowService, ImportWorkflowService, ListTaskService, + BlobAuditWorkflowService, RebuildFTSIndexWorkflowService, SetDownloadFilesService, TaskListResponse, @@ -2287,6 +2288,23 @@ export function sendFullTextSearch(query: string, offset?: number): ThunkRespons }; } +export function sendBlobAuditTask(req: BlobAuditWorkflowService): ThunkResponse { + return async (dispatch, _getState) => { + return await dispatch( + send( + "/workflow/blobAudit", + { + data: req, + method: "POST", + }, + { + ...defaultOpts, + }, + ), + ); + }; +} + export function sendRebuildFTSIndex(req: RebuildFTSIndexWorkflowService): ThunkResponse { return async (dispatch, _getState) => { return await dispatch( diff --git a/frontend/src/api/workflow.ts b/frontend/src/api/workflow.ts index 521c4c1c..601c3908 100644 --- a/frontend/src/api/workflow.ts +++ b/frontend/src/api/workflow.ts @@ -142,6 +142,7 @@ export enum TaskType { full_text_change_owner = "full_text_change_owner", full_text_delete = "full_text_delete", full_text_rebuild = "full_text_rebuild", + blob_audit = "blob_audit", } export enum TaskStatus { @@ -175,3 +176,8 @@ export interface SetDownloadFilesService { export interface RebuildFTSIndexWorkflowService { filtered_storage_policy?: number[]; } + +export interface BlobAuditWorkflowService { + policy_id: number; + delete?: boolean; +} diff --git a/frontend/src/component/Admin/StoragePolicy/BlobAuditDialog.tsx b/frontend/src/component/Admin/StoragePolicy/BlobAuditDialog.tsx new file mode 100644 index 00000000..a6c1bb21 --- /dev/null +++ b/frontend/src/component/Admin/StoragePolicy/BlobAuditDialog.tsx @@ -0,0 +1,84 @@ +import { + Button, + Checkbox, + Dialog, + DialogActions, + DialogContent, + DialogContentText, + DialogTitle, + FormControlLabel, +} from "@mui/material"; +import { useSnackbar } from "notistack"; +import { useEffect, useState } from "react"; +import { useTranslation } from "react-i18next"; +import { sendBlobAuditTask } from "../../../api/api"; +import { StoragePolicy } from "../../../api/dashboard"; +import { useAppDispatch } from "../../../redux/hooks"; +import { DefaultCloseAction } from "../../Common/Snackbar/snackbar"; + +export interface BlobAuditDialogProps { + open: boolean; + onClose: () => void; + policy?: StoragePolicy; +} + +const BlobAuditDialog = ({ open, onClose, policy }: BlobAuditDialogProps) => { + const { t } = useTranslation("dashboard"); + const dispatch = useAppDispatch(); + const { enqueueSnackbar } = useSnackbar(); + const [deleteOrphans, setDeleteOrphans] = useState(false); + const [loading, setLoading] = useState(false); + + useEffect(() => { + if (open) { + setDeleteOrphans(false); + } + }, [open]); + + const onSubmit = () => { + if (!policy?.id) { + return; + } + setLoading(true); + dispatch(sendBlobAuditTask({ policy_id: policy.id, delete: deleteOrphans })) + .then(() => { + enqueueSnackbar(t("policy.blobAuditSubmitted"), { variant: "success", action: DefaultCloseAction }); + onClose(); + }) + .finally(() => { + setLoading(false); + }); + }; + + return ( + + {t("policy.blobAudit", { name: policy?.name ?? "" })} + + {t("policy.blobAuditDes")} + setDeleteOrphans(e.target.checked)} + /> + } + label={t("policy.blobAuditDelete")} + /> + {deleteOrphans && ( + + {t("policy.blobAuditDeleteWarning")} + + )} + + + + + + + ); +}; + +export default BlobAuditDialog; diff --git a/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx b/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx index c59f9b57..7f03a4f8 100644 --- a/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx +++ b/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx @@ -1,4 +1,5 @@ -import { Box, Divider, IconButton, Link, Skeleton, Typography } from "@mui/material"; +import { FindInPage } from "@mui/icons-material"; +import { Box, Divider, IconButton, Link, Skeleton, Tooltip, Typography } from "@mui/material"; import Grid from "@mui/material/Grid2"; import { useCallback, useState } from "react"; import { useTranslation } from "react-i18next"; @@ -12,6 +13,7 @@ import { NoWrapBox, SquareChip } from "../../Common/StyledComponents"; import Delete from "../../Icons/Delete"; import Info from "../../Icons/Info"; import { BorderedCardClickableBaImg } from "../Common/AdminCard"; +import BlobAuditDialog from "./BlobAuditDialog"; import { PolicyPropsMap } from "./StoragePolicySetting"; export interface StoragePolicyCardProps { @@ -26,8 +28,14 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp const [detail, setDetail] = useState(undefined); const [deleteLoading, setDeleteLoading] = useState(false); const [detailLoading, setDetailLoading] = useState(false); + const [auditOpen, setAuditOpen] = useState(false); const navigate = useNavigate(); + const handleAuditOpen = useCallback((e: React.MouseEvent) => { + e.stopPropagation(); + setAuditOpen(true); + }, []); + const loadDetail = useCallback( (e: React.MouseEvent) => { e.preventDefault(); @@ -169,11 +177,19 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp )} - - - + + + + + + + + + + + setAuditOpen(false)} policy={policy} /> ); }; diff --git a/frontend/src/component/Admin/Task/TaskContent.tsx b/frontend/src/component/Admin/Task/TaskContent.tsx index 5819b1cf..7612ccd5 100644 --- a/frontend/src/component/Admin/Task/TaskContent.tsx +++ b/frontend/src/component/Admin/Task/TaskContent.tsx @@ -133,6 +133,11 @@ export const TaskContent = memo(({ task, openEntity, openFile }: TaskContentProp return t("task.fullTextDelete", { count: privateState?.file_ids?.length ?? 0, }); + case TaskType.blob_audit: + return t("task.blobAudit", { + policyID: privateState?.policy_id ?? 0, + count: privateState?.orphan_count ?? 0, + }); default: return ""; } diff --git a/pkg/filemanager/workflows/blob_audit.go b/pkg/filemanager/workflows/blob_audit.go new file mode 100644 index 00000000..e1a54541 --- /dev/null +++ b/pkg/filemanager/workflows/blob_audit.go @@ -0,0 +1,301 @@ +package workflows + +import ( + "context" + "encoding/json" + "fmt" + "path" + "strings" + "sync/atomic" + + "github.com/cloudreve/Cloudreve/v4/application/dependency" + "github.com/cloudreve/Cloudreve/v4/ent" + "github.com/cloudreve/Cloudreve/v4/ent/storagepolicy" + "github.com/cloudreve/Cloudreve/v4/ent/task" + "github.com/cloudreve/Cloudreve/v4/inventory" + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/manager" + "github.com/cloudreve/Cloudreve/v4/pkg/hashid" + "github.com/cloudreve/Cloudreve/v4/pkg/logging" + "github.com/cloudreve/Cloudreve/v4/pkg/queue" +) + +type ( + BlobAuditTask struct { + *queue.DBTask + state *BlobAuditTaskState + l logging.Logger + progress queue.Progresses + } + BlobAuditTaskPhase string + BlobAuditTaskState struct { + Phase BlobAuditTaskPhase `json:"phase"` + PolicyID int `json:"policy_id"` + Delete bool `json:"delete"` + Scanned int `json:"scanned"` + Orphans []string `json:"orphans"` + OrphanCount int `json:"orphan_count"` + Deleted int `json:"deleted"` + Failed int `json:"failed"` + Truncated bool `json:"truncated"` + } +) + +const ( + BlobAuditPhaseScan BlobAuditTaskPhase = "scan" + BlobAuditPhaseDelete BlobAuditTaskPhase = "delete" + + BlobAuditDeleteBatch = 200 + BlobAuditMaxOrphanList = 5000 + + ProgressTypeBlobAudit = "blob_audit" + SummaryKeyOrphanCount = "orphan_count" + SummaryKeyDeleted = "deleted" + SummaryKeyScannedObject = "scanned" +) + +func init() { + queue.RegisterResumableTaskFactory(queue.BlobAuditTaskType, NewBlobAuditTaskFromModel) +} + +// NewBlobAuditTask creates a task that lists all physical objects under a +// storage policy and diffs them against the entities table, reporting (and +// optionally deleting) blobs no entity references (upstream #3395). +func NewBlobAuditTask(ctx context.Context, u *ent.User, policyID int, delete bool) (queue.Task, error) { + state := &BlobAuditTaskState{ + Phase: BlobAuditPhaseScan, + PolicyID: policyID, + Delete: delete, + } + stateBytes, err := json.Marshal(state) + if err != nil { + return nil, fmt.Errorf("failed to marshal state: %w", err) + } + + return &BlobAuditTask{ + DBTask: &queue.DBTask{ + Task: &ent.Task{ + Type: queue.BlobAuditTaskType, + CorrelationID: logging.CorrelationID(ctx), + PrivateState: string(stateBytes), + PublicState: &types.TaskPublicState{}, + }, + DirectOwner: u, + }, + }, nil +} + +func NewBlobAuditTaskFromModel(t *ent.Task) queue.Task { + return &BlobAuditTask{ + DBTask: &queue.DBTask{ + Task: t, + }, + } +} + +func (m *BlobAuditTask) Do(ctx context.Context) (task.Status, error) { + dep := dependency.FromContext(ctx) + m.l = dep.Logger() + + m.Lock() + if m.progress == nil { + m.progress = make(queue.Progresses) + } + m.progress[ProgressTypeBlobAudit] = &queue.Progress{} + m.Unlock() + + state := &BlobAuditTaskState{} + if err := json.Unmarshal([]byte(m.State()), state); err != nil { + return task.StatusError, fmt.Errorf("failed to unmarshal state: %s (%w)", err, queue.CriticalErr) + } + m.state = state + + var ( + next = task.StatusCompleted + err error + ) + switch m.state.Phase { + case BlobAuditPhaseScan, "": + next, err = m.scan(ctx, dep) + case BlobAuditPhaseDelete: + next, err = m.deleteOrphans(ctx, dep) + default: + next, err = task.StatusError, fmt.Errorf("unknown phase %q: %w", m.state.Phase, queue.CriticalErr) + } + + newStateStr, marshalErr := json.Marshal(m.state) + if marshalErr != nil { + return task.StatusError, fmt.Errorf("failed to marshal state: %w", marshalErr) + } + + m.Lock() + m.Task.PrivateState = string(newStateStr) + m.Unlock() + return next, err +} + +// scan lists all physical objects under the policy's blob prefix and diffs +// them against the entity sources recorded for that policy. +func (m *BlobAuditTask) scan(ctx context.Context, dep dependency.Dep) (task.Status, error) { + policy, err := dep.StoragePolicyClient().GetPolicyByID(ctx, m.state.PolicyID) + if err != nil { + return task.StatusError, fmt.Errorf("failed to get storage policy %d: %w", m.state.PolicyID, err) + } + if policy.Status == storagepolicy.StatusSuspended { + return task.StatusError, fmt.Errorf("storage policy %d is suspended: %w", m.state.PolicyID, queue.CriticalErr) + } + + fm := manager.NewFileManager(dep, nil) + defer fm.Recycle() + d, err := fm.GetStorageDriver(ctx, policy) + if err != nil { + return task.StatusError, fmt.Errorf("failed to get storage driver: %w", err) + } + + // Literal prefix of the dir name rule (e.g. "uploads/{uid}/{path}" -> + // "uploads") scopes the listing so non-blob areas (temp, avatars) are + // never touched. + base := literalDirPrefix(policy.DirNameRule) + m.l.Info("Blob audit on policy %d: listing objects under %q", policy.ID, base) + + var objects []fs.PhysicalObject + objects, err = d.List(ctx, base, func(i int) { + atomic.AddInt64(&m.progress[ProgressTypeBlobAudit].Current, int64(i)) + }, true) + if err != nil { + return task.StatusError, fmt.Errorf("failed to list objects on policy %d: %w", m.state.PolicyID, err) + } + m.state.Scanned = len(objects) + + // Full entity source set for this policy, paginated. + known := make(map[string]struct{}) + page := 0 + for { + res, err := dep.FileClient().ListEntities(ctx, &inventory.ListEntityParameters{ + PaginationArgs: &inventory.PaginationArgs{Page: page, PageSize: 1000}, + StoragePolicyID: policy.ID, + }) + if err != nil { + return task.StatusError, fmt.Errorf("failed to list entities: %w", err) + } + for _, e := range res.Entities { + known[e.Source] = struct{}{} + } + if len(res.Entities) < 1000 { + break + } + page++ + } + + orphanSet := make(map[string]struct{}) + for _, obj := range objects { + if obj.IsDir { + continue + } + key := obj.RelativePath + if base != "" { + key = path.Join(base, obj.RelativePath) + } + if _, ok := known[key]; ok { + continue + } + if _, dup := orphanSet[key]; dup { + continue + } + orphanSet[key] = struct{}{} + if len(m.state.Orphans) < BlobAuditMaxOrphanList { + m.state.Orphans = append(m.state.Orphans, key) + } else { + m.state.Truncated = true + } + } + + m.l.Info("Blob audit on policy %d: %d objects scanned, %d orphans found.", policy.ID, m.state.Scanned, len(orphanSet)) + m.state.OrphanCount = len(orphanSet) + + if m.state.Delete && len(m.state.Orphans) > 0 { + m.state.Phase = BlobAuditPhaseDelete + m.ResumeAfter(0) + return task.StatusSuspending, nil + } + return task.StatusCompleted, nil +} + +// deleteOrphans removes the recorded orphan blobs in batches. +func (m *BlobAuditTask) deleteOrphans(ctx context.Context, dep dependency.Dep) (task.Status, error) { + if len(m.state.Orphans) == 0 { + return task.StatusCompleted, nil + } + + policy, err := dep.StoragePolicyClient().GetPolicyByID(ctx, m.state.PolicyID) + if err != nil { + return task.StatusError, fmt.Errorf("failed to get storage policy %d: %w", m.state.PolicyID, err) + } + + fm := manager.NewFileManager(dep, nil) + defer fm.Recycle() + d, err := fm.GetStorageDriver(ctx, policy) + if err != nil { + return task.StatusError, fmt.Errorf("failed to get storage driver: %w", err) + } + + batch := m.state.Orphans + if len(batch) > BlobAuditDeleteBatch { + batch = batch[:BlobAuditDeleteBatch] + } + + failed, err := d.Delete(ctx, batch...) + if err != nil { + m.l.Warning("Blob audit delete partially failed on policy %d: %s", policy.ID, err) + } + m.state.Deleted += len(batch) - len(failed) + m.state.Failed += len(failed) + // Failed paths are dropped from the list — a re-run of the audit flags + // them again if they are still orphaned. + m.state.Orphans = m.state.Orphans[len(batch):] + + atomic.StoreInt64(&m.progress[ProgressTypeBlobAudit].Current, int64(m.state.Deleted)) + + if len(m.state.Orphans) > 0 { + m.ResumeAfter(0) + return task.StatusSuspending, nil + } + return task.StatusCompleted, nil +} + +// literalDirPrefix returns the static path prefix before the first magic +// variable, e.g. "uploads/{uid}/{path}" -> "uploads". An empty result means +// the rule starts with a variable and the policy root cannot be narrowed. +func literalDirPrefix(rule string) string { + if i := strings.Index(rule, "{"); i >= 0 { + rule = rule[:i] + } + return strings.Trim(rule, "/") +} + +func (m *BlobAuditTask) Progress(ctx context.Context) queue.Progresses { + m.Lock() + defer m.Unlock() + return m.progress +} + +func (m *BlobAuditTask) Summarize(hasher hashid.Encoder) *queue.Summary { + if m.state == nil { + if err := json.Unmarshal([]byte(m.State()), &m.state); err != nil { + return nil + } + } + + return &queue.Summary{ + Phase: string(m.state.Phase), + Props: map[string]any{ + SummaryKeyOrphanCount: m.state.OrphanCount, + SummaryKeyScannedObject: m.state.Scanned, + SummaryKeyDeleted: m.state.Deleted, + SummaryKeyFailed: m.state.Failed, + "truncated": m.state.Truncated, + "policy_id": m.state.PolicyID, + }, + } +} diff --git a/pkg/filemanager/workflows/blob_audit_test.go b/pkg/filemanager/workflows/blob_audit_test.go new file mode 100644 index 00000000..7e29164c --- /dev/null +++ b/pkg/filemanager/workflows/blob_audit_test.go @@ -0,0 +1,23 @@ +package workflows + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestLiteralDirPrefix(t *testing.T) { + cases := []struct { + rule string + want string + }{ + {"uploads/{uid}/{path}", "uploads"}, + {"uploads/{uid}", "uploads"}, + {"{uid}/{path}", ""}, + {"blob/store/{uid}", "blob/store"}, + {"", ""}, + } + for _, c := range cases { + require.Equal(t, c.want, literalDirPrefix(c.rule), "rule %q", c.rule) + } +} diff --git a/pkg/queue/task.go b/pkg/queue/task.go index 53d27fd3..a648efbe 100644 --- a/pkg/queue/task.go +++ b/pkg/queue/task.go @@ -104,6 +104,7 @@ const ( RelocateTaskType = "relocate" RemoteDownloadTaskType = "remote_download" ImportTaskType = "import" + BlobAuditTaskType = "blob_audit" FullTextIndexTaskType = "full_text_index" FullTextCopyTaskType = "full_text_copy" diff --git a/routers/controllers/file.go b/routers/controllers/file.go index a956c6b2..4c08e5a1 100644 --- a/routers/controllers/file.go +++ b/routers/controllers/file.go @@ -83,6 +83,21 @@ func RebuildFTSIndex(c *gin.Context) { }) } +// BlobAudit creates a blob-vs-database audit task +func BlobAudit(c *gin.Context) { + service := ParametersFromContext[*explorer.BlobAuditWorkflowService](c, explorer.BlobAuditParamCtx{}) + resp, err := service.CreateBlobAuditTask(c) + if err != nil { + c.JSON(200, serializer.Err(c, err)) + c.Abort() + return + } + + c.JSON(200, serializer.Response{ + Data: resp, + }) +} + // ExtractArchive creates extract archive task func ExtractArchive(c *gin.Context) { service := ParametersFromContext[*explorer.ArchiveWorkflowService](c, explorer.CreateArchiveParamCtx{}) diff --git a/routers/router.go b/routers/router.go index 2776a7da..8aef1436 100644 --- a/routers/router.go +++ b/routers/router.go @@ -814,6 +814,13 @@ func initMasterRouter(dep dependency.Dep) *gin.Engine { controllers.FromJSON[explorer.RebuildFTSIndexWorkflowService](explorer.CreateRebuildFTSIndexParamCtx{}), controllers.RebuildFTSIndex, ) + // Create task to audit physical blobs against the entities table + wf.POST("blobAudit", + middleware.IsAdmin(), + middleware.RequiredScopes(types.ScopeWorkflowWrite, types.ScopeAdminWrite), + controllers.FromJSON[explorer.BlobAuditWorkflowService](explorer.BlobAuditParamCtx{}), + controllers.BlobAudit, + ) // 取得文件外链 source := file.Group("source") diff --git a/service/explorer/workflows.go b/service/explorer/workflows.go index 4b385a6d..4f8da763 100644 --- a/service/explorer/workflows.go +++ b/service/explorer/workflows.go @@ -645,3 +645,37 @@ func (service *RebuildFTSIndexWorkflowService) CreateRebuildFTSIndexTask(c *gin. return BuildTaskResponse(t, nil, hasher), nil } + +type ( + BlobAuditWorkflowService struct { + PolicyID int `json:"policy_id" binding:"required,min=1"` + Delete bool `json:"delete"` + } + BlobAuditParamCtx struct{} +) + +// CreateBlobAuditTask queues a blob-vs-database audit for one storage policy. +func (service *BlobAuditWorkflowService) CreateBlobAuditTask(c *gin.Context) (*TaskResponse, error) { + dep := dependency.FromContext(c) + user := inventory.UserFromContext(c) + hasher := dep.HashIDEncoder() + + if !user.Edges.Group.Permissions.Enabled(int(types.GroupPermissionIsAdmin)) { + return nil, serializer.NewError(serializer.CodeGroupNotAllowed, "Only admin can run a blob audit", nil) + } + + if _, err := dep.StoragePolicyClient().GetPolicyByID(c, service.PolicyID); err != nil { + return nil, serializer.NewError(serializer.CodeNotFound, "Storage policy not found", err) + } + + t, err := workflows.NewBlobAuditTask(c, user, service.PolicyID, service.Delete) + if err != nil { + return nil, serializer.NewError(serializer.CodeCreateTaskError, "Failed to create task", err) + } + + if err := dep.IoIntenseQueue(c).QueueTask(c, t); err != nil { + return nil, serializer.NewError(serializer.CodeCreateTaskError, "Failed to queue task", err) + } + + return BuildTaskResponse(t, nil, hasher), nil +}