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>
pull/3582/head
Tomas Dvorak 2 weeks ago
parent 5ecfd3e9f1
commit 09031da204

@ -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}}</0>",
"fullTextChangeOwner": "Change index owner for file <0>#{{fileID}}</0>",
"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",

@ -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}}</0> 的全文索引",
"fullTextChangeOwner": "变更文件 <0>#{{fileID}}</0> 的索引所有者",
"fullTextDelete": "删除 {{count}} 个文件的全文索引",
"blobAudit": "审计存储策略 #{{policyID}} 的 Blob({{count}} 个孤儿对象)",
"type": "类型",
"node": "处理节点",
"createdBy": "创建者",

@ -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<TaskResponse> {
return async (dispatch, _getState) => {
return await dispatch(
send(
"/workflow/blobAudit",
{
data: req,
method: "POST",
},
{
...defaultOpts,
},
),
);
};
}
export function sendRebuildFTSIndex(req: RebuildFTSIndexWorkflowService): ThunkResponse<TaskResponse> {
return async (dispatch, _getState) => {
return await dispatch(

@ -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;
}

@ -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 (
<Dialog open={open} onClose={onClose} maxWidth="xs" fullWidth>
<DialogTitle>{t("policy.blobAudit", { name: policy?.name ?? "" })}</DialogTitle>
<DialogContent>
<DialogContentText sx={{ mb: 2 }}>{t("policy.blobAuditDes")}</DialogContentText>
<FormControlLabel
control={
<Checkbox
size="small"
checked={deleteOrphans}
onChange={(e) => setDeleteOrphans(e.target.checked)}
/>
}
label={t("policy.blobAuditDelete")}
/>
{deleteOrphans && (
<DialogContentText color="error" variant="body2">
{t("policy.blobAuditDeleteWarning")}
</DialogContentText>
)}
</DialogContent>
<DialogActions>
<Button onClick={onClose}>{t("common:cancel")}</Button>
<Button variant="contained" color={deleteOrphans ? "error" : "primary"} onClick={onSubmit} disabled={loading}>
{t("policy.blobAuditStart")}
</Button>
</DialogActions>
</Dialog>
);
};
export default BlobAuditDialog;

@ -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<StoragePolicy | undefined>(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<HTMLButtonElement>) => {
e.stopPropagation();
setAuditOpen(true);
}, []);
const loadDetail = useCallback(
(e: React.MouseEvent<HTMLAnchorElement>) => {
e.preventDefault();
@ -169,11 +177,19 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp
)}
</Typography>
<Box sx={{ display: "flex", alignItems: "center" }}>
<Tooltip title={t("policy.blobAudit", { name: policy?.name ?? "" })}>
<IconButton size="small" onClick={handleAuditOpen} disabled={deleteLoading}>
<FindInPage fontSize="small" />
</IconButton>
</Tooltip>
<IconButton size="small" onClick={handleDelete}>
<Delete fontSize="small" />
</IconButton>
</Box>
</Box>
</BorderedCardClickableBaImg>
<BlobAuditDialog open={auditOpen} onClose={() => setAuditOpen(false)} policy={policy} />
</Grid>
);
};

@ -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 "";
}

@ -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,
},
}
}

@ -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)
}
}

@ -104,6 +104,7 @@ const (
RelocateTaskType = "relocate"
RemoteDownloadTaskType = "remote_download"
ImportTaskType = "import"
BlobAuditTaskType = "blob_audit"
FullTextIndexTaskType = "full_text_index"
FullTextCopyTaskType = "full_text_copy"

@ -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{})

@ -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")

@ -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
}

Loading…
Cancel
Save