diff --git a/frontend/public/locales/en-US/dashboard.json b/frontend/public/locales/en-US/dashboard.json index 60867d21..ffc04805 100644 --- a/frontend/public/locales/en-US/dashboard.json +++ b/frontend/public/locales/en-US/dashboard.json @@ -1518,7 +1518,14 @@ "confirmBatchDelete": "Are you sure you want to delete {{num}} Blobs?", "deleteXEntities": "Delete {{num}} Blobs", "forceDelete": "Force delete", - "forceDeleteDes": "Whether to delete the Blob record regardless of whether the physical file is deleted." + "forceDeleteDes": "Whether to delete the Blob record regardless of whether the physical file is deleted.", + "relocateXEntities": "Relocate {{num}} Blobs", + "relocateTitle": "Relocate Blobs", + "relocateEntitiesDes": "Move the data of {{num}} selected Blobs to another storage policy. The task runs in the background and can be resumed.", + "relocatePolicyDes": "Move all Blobs stored on \"{{name}}\" to another storage policy. The task runs in the background and can be resumed.", + "relocatePolicy": "Migrate all Blobs on {{name}} to another policy", + "relocateStart": "Start relocation", + "relocateSubmitted": "Relocation task submitted" }, "event": { "cleanup": "Cleanup", diff --git a/frontend/public/locales/zh-CN/dashboard.json b/frontend/public/locales/zh-CN/dashboard.json index ddfc2272..98e9982f 100644 --- a/frontend/public/locales/zh-CN/dashboard.json +++ b/frontend/public/locales/zh-CN/dashboard.json @@ -1518,7 +1518,14 @@ "confirmBatchDelete": "确认要删除 {{num}} 个 Blob?", "deleteXEntities": "删除 {{num}} 个 Blob", "forceDelete": "强制删除", - "forceDeleteDes": "无论物理文件是否删除成功,都会删除 Blob 记录。" + "forceDeleteDes": "无论物理文件是否删除成功,都会删除 Blob 记录。", + "relocateXEntities": "迁移 {{num}} 个 Blob", + "relocateTitle": "迁移 Blob", + "relocateEntitiesDes": "将选中的 {{num}} 个 Blob 的数据迁移到另一个存储策略。任务将在后台运行,可断点续传。", + "relocatePolicyDes": "将存储在“{{name}}”上的所有 Blob 迁移到另一个存储策略。任务将在后台运行,可断点续传。", + "relocatePolicy": "将 {{name}} 上的所有 Blob 迁移到其他策略", + "relocateStart": "开始迁移", + "relocateSubmitted": "迁移任务已提交" }, "event": { "cleanup": "清理", diff --git a/frontend/src/api/api.ts b/frontend/src/api/api.ts index 8ab5f525..2c820e65 100644 --- a/frontend/src/api/api.ts +++ b/frontend/src/api/api.ts @@ -108,6 +108,7 @@ import { ListTaskService, BlobAuditWorkflowService, RebuildFTSIndexWorkflowService, + RelocateEntityService, SetDownloadFilesService, TaskListResponse, TaskProgresses, @@ -2339,6 +2340,23 @@ export function sendBlobAuditTask(req: BlobAuditWorkflowService): ThunkResponse< }; } +export function relocateEntities(req: RelocateEntityService): ThunkResponse { + return async (dispatch, _getState) => { + return await dispatch( + send( + "/admin/entity/relocate", + { + 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 5fd805d9..e28ec1aa 100644 --- a/frontend/src/api/workflow.ts +++ b/frontend/src/api/workflow.ts @@ -182,3 +182,9 @@ export interface BlobAuditWorkflowService { policy_id: number; delete?: boolean; } + +export interface RelocateEntityService { + entity_ids?: number[]; + src_policy_id?: number; + dst_policy_id: number; +} diff --git a/frontend/src/component/Admin/Entity/EntitySetting.tsx b/frontend/src/component/Admin/Entity/EntitySetting.tsx index a6f98088..550a7c30 100644 --- a/frontend/src/component/Admin/Entity/EntitySetting.tsx +++ b/frontend/src/component/Admin/Entity/EntitySetting.tsx @@ -1,4 +1,4 @@ -import { Delete } from "@mui/icons-material"; +import { Delete, DriveFileMove } from "@mui/icons-material"; import { Badge, Box, @@ -36,6 +36,7 @@ import EntityDeleteDialog from "./EntityDeleteDialog"; import EntityDialog from "./EntityDialog/EntityDialog"; import EntityFilterPopover from "./EntityFilterPopover"; import EntityRow from "./EntityRow"; +import RelocateDialog from "./RelocateDialog"; export const StoragePolicyQuery = "storage_policy"; export const UserQuery = "user"; export const TypeQuery = "type"; @@ -73,6 +74,7 @@ const EntitySetting = () => { const [entityDialogID, setEntityDialogID] = useState(undefined); const [deleteDialogOpen, setDeleteDialogOpen] = useState(false); const [deleteDialogID, setDeleteDialogID] = useState(undefined); + const [relocateDialogOpen, setRelocateDialogOpen] = useState(false); const pageInt = parseInt(page) ?? 1; const pageSizeInt = parseInt(pageSize) ?? 10; @@ -198,6 +200,11 @@ const EntitySetting = () => { entityID={deleteDialogID} onDelete={fetchEntities} /> + setRelocateDialogOpen(false)} + entityIDs={Array.from(selected)} + /> @@ -227,6 +234,14 @@ const EntitySetting = () => { {selected.length > 0 && !isMobile && ( <> + @@ -235,6 +250,14 @@ const EntitySetting = () => { {isMobile && selected.length > 0 && ( + diff --git a/frontend/src/component/Admin/Entity/RelocateDialog.tsx b/frontend/src/component/Admin/Entity/RelocateDialog.tsx new file mode 100644 index 00000000..d3e75092 --- /dev/null +++ b/frontend/src/component/Admin/Entity/RelocateDialog.tsx @@ -0,0 +1,73 @@ +import { Button, Dialog, DialogActions, DialogContent, DialogContentText, DialogTitle } from "@mui/material"; +import { useSnackbar } from "notistack"; +import { useEffect, useState } from "react"; +import { useTranslation } from "react-i18next"; +import { relocateEntities } from "../../../api/api"; +import { useAppDispatch } from "../../../redux/hooks"; +import { DefaultCloseAction } from "../../Common/Snackbar/snackbar"; +import SinglePolicySelectionInput from "../Common/SinglePolicySelectionInput"; + +export interface RelocateDialogProps { + open: boolean; + onClose: () => void; + entityIDs?: number[]; + srcPolicyID?: number; + srcPolicyName?: string; +} + +const RelocateDialog = ({ open, onClose, entityIDs, srcPolicyID, srcPolicyName }: RelocateDialogProps) => { + const { t } = useTranslation("dashboard"); + const dispatch = useAppDispatch(); + const { enqueueSnackbar } = useSnackbar(); + const [dstPolicyID, setDstPolicyID] = useState(undefined); + const [loading, setLoading] = useState(false); + + useEffect(() => { + if (open) { + setDstPolicyID(undefined); + } + }, [open]); + + const onSubmit = () => { + if (!dstPolicyID) { + return; + } + setLoading(true); + dispatch( + relocateEntities({ + entity_ids: entityIDs, + src_policy_id: srcPolicyID, + dst_policy_id: dstPolicyID, + }), + ) + .then(() => { + enqueueSnackbar(t("entity.relocateSubmitted"), { variant: "success", action: DefaultCloseAction }); + onClose(); + }) + .finally(() => { + setLoading(false); + }); + }; + + return ( + + {t("entity.relocateTitle")} + + + {srcPolicyID + ? t("entity.relocatePolicyDes", { name: srcPolicyName ?? "" }) + : t("entity.relocateEntitiesDes", { num: entityIDs?.length ?? 0 })} + + + + + + + + + ); +}; + +export default RelocateDialog; diff --git a/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx b/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx index 7f03a4f8..90b9190b 100644 --- a/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx +++ b/frontend/src/component/Admin/StoragePolicy/StoragePolicyCard.tsx @@ -1,4 +1,4 @@ -import { FindInPage } from "@mui/icons-material"; +import { DriveFileMove, 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"; @@ -13,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 RelocateDialog from "../Entity/RelocateDialog"; import BlobAuditDialog from "./BlobAuditDialog"; import { PolicyPropsMap } from "./StoragePolicySetting"; @@ -29,6 +30,7 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp const [deleteLoading, setDeleteLoading] = useState(false); const [detailLoading, setDetailLoading] = useState(false); const [auditOpen, setAuditOpen] = useState(false); + const [relocateOpen, setRelocateOpen] = useState(false); const navigate = useNavigate(); const handleAuditOpen = useCallback((e: React.MouseEvent) => { @@ -36,6 +38,11 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp setAuditOpen(true); }, []); + const handleRelocateOpen = useCallback((e: React.MouseEvent) => { + e.stopPropagation(); + setRelocateOpen(true); + }, []); + const loadDetail = useCallback( (e: React.MouseEvent) => { e.preventDefault(); @@ -178,6 +185,11 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp + + + + + @@ -190,6 +202,12 @@ const StoragePolicyCard = ({ policy, onRefresh, loading }: StoragePolicyCardProp setAuditOpen(false)} policy={policy} /> + setRelocateOpen(false)} + srcPolicyID={policy?.id} + srcPolicyName={policy?.name} + /> ); }; diff --git a/inventory/file.go b/inventory/file.go index 5f1abfde..98d45b77 100644 --- a/inventory/file.go +++ b/inventory/file.go @@ -204,6 +204,9 @@ type FileClient interface { RemoveEntitiesByID(ctx context.Context, ids ...int) (map[int]int64, error) // UpdateEntityProps persists mutated EntityProps back to the given entities. UpdateEntityProps(ctx context.Context, entities ...*ent.Entity) error + // RelocateEntity moves an entity's blob record to a different storage policy: + // updates the blob path, the owning policy, and props (e.g. encryption metadata). + RelocateEntity(ctx context.Context, entity *ent.Entity, newSource string, dstPolicyID int) error // CapEntities caps the number of entities of a given file. The oldest entities will be unlinked // if entity count exceed limit. CapEntities(ctx context.Context, file *ent.File, owner *ent.User, max int, entityType types.EntityType) (StorageDiff, error) @@ -475,6 +478,18 @@ func (f *fileClient) RemoveEntitiesByID(ctx context.Context, ids ...int) (map[in return storageReduced, nil } +func (f *fileClient) RelocateEntity(ctx context.Context, entity *ent.Entity, newSource string, dstPolicyID int) error { + err := f.client.Entity.UpdateOne(entity). + SetSource(newSource). + SetStoragePolicyEntities(dstPolicyID). + SetProps(entity.Props). + Exec(ctx) + if err != nil { + return fmt.Errorf("failed to relocate entity %d to policy %d: %w", entity.ID, dstPolicyID, err) + } + return nil +} + func (f *fileClient) UpdateEntityProps(ctx context.Context, entities ...*ent.Entity) error { for _, e := range entities { if err := f.client.Entity.UpdateOne(e).SetProps(e.Props).Exec(ctx); err != nil { diff --git a/pkg/filemanager/workflows/relocate.go b/pkg/filemanager/workflows/relocate.go new file mode 100644 index 00000000..32e64ef9 --- /dev/null +++ b/pkg/filemanager/workflows/relocate.go @@ -0,0 +1,451 @@ +package workflows + +import ( + "context" + "encoding/json" + "fmt" + "io" + "path" + "sort" + "strconv" + "sync/atomic" + "time" + + "github.com/cloudreve/Cloudreve/v4/application/dependency" + "github.com/cloudreve/Cloudreve/v4/ent" + "github.com/cloudreve/Cloudreve/v4/ent/entity" + "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 ( + RelocateTask struct { + *queue.DBTask + + l logging.Logger + state *RelocateTaskState + progress queue.Progresses + } + RelocateTaskPhase string + RelocateTaskState struct { + // EntityIDs is an explicit set of entities to relocate, sorted ascending + // at creation so Cursor resume works uniformly. + EntityIDs []int `json:"entity_ids,omitempty"` + // SrcPolicyID relocates every entity still on this policy instead of an + // explicit list. + SrcPolicyID int `json:"src_policy_id,omitempty"` + DstPolicyID int `json:"dst_policy_id"` + // Cursor is the ID of the last successfully processed entity. + Cursor int `json:"cursor,omitempty"` + // Failed maps entity ID to the error that prevented its relocation. + Failed map[string]string `json:"failed,omitempty"` + NodeState `json:",inline"` + Phase RelocateTaskPhase `json:"phase,omitempty"` + } +) + +const ( + RelocatePhaseNotStarted RelocateTaskPhase = "" + RelocatePhaseRelocating RelocateTaskPhase = "relocating" + + ProgressTypeRelocateCount = "relocate_count" + ProgressTypeRelocateSize = "relocate_size" + + relocatePageSize = 50 + + SummaryKeyDstPolicy = "dst_policy" + SummaryKeySrcPolicy = "src_policy" + SummaryKeyEntities = "entities" + SummaryKeyRelocateFailed = "relocate_failed" +) + +func init() { + queue.RegisterResumableTaskFactory(queue.RelocateTaskType, NewRelocateTaskFromModel) +} + +// NewRelocateTask creates a RelocateTask moving the given entities to dstPolicyID. +func NewRelocateTask(ctx context.Context, entityIDs []int, dstPolicyID int) (queue.Task, error) { + ids := append([]int(nil), entityIDs...) + sort.Ints(ids) + return newRelocateTask(ctx, &RelocateTaskState{ + EntityIDs: ids, + DstPolicyID: dstPolicyID, + NodeState: NodeState{}, + }) +} + +// NewRelocatePolicyTask creates a RelocateTask moving every entity on +// srcPolicyID to dstPolicyID. +func NewRelocatePolicyTask(ctx context.Context, srcPolicyID, dstPolicyID int) (queue.Task, error) { + return newRelocateTask(ctx, &RelocateTaskState{ + SrcPolicyID: srcPolicyID, + DstPolicyID: dstPolicyID, + NodeState: NodeState{}, + }) +} + +func newRelocateTask(ctx context.Context, state *RelocateTaskState) (queue.Task, error) { + stateBytes, err := json.Marshal(state) + if err != nil { + return nil, fmt.Errorf("failed to marshal state: %w", err) + } + + return &RelocateTask{ + DBTask: &queue.DBTask{ + Task: &ent.Task{ + Type: queue.RelocateTaskType, + CorrelationID: logging.CorrelationID(ctx), + PrivateState: string(stateBytes), + PublicState: &types.TaskPublicState{}, + }, + DirectOwner: inventory.UserFromContext(ctx), + }, + }, nil +} + +func NewRelocateTaskFromModel(task *ent.Task) queue.Task { + return &RelocateTask{ + DBTask: &queue.DBTask{ + Task: task, + }, + } +} + +func (m *RelocateTask) 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.Unlock() + + state := &RelocateTaskState{} + if err := json.Unmarshal([]byte(m.State()), state); err != nil { + return task.StatusError, fmt.Errorf("failed to unmarshal state: %w", err) + } + m.state = state + + next, err := m.relocate(ctx, dep) + + 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 +} + +func (m *RelocateTask) relocate(ctx context.Context, dep dependency.Dep) (task.Status, error) { + if m.state.SrcPolicyID != 0 && m.state.SrcPolicyID == m.state.DstPolicyID { + return task.StatusError, fmt.Errorf("source and destination policies are identical (%w)", queue.CriticalErr) + } + + policyClient := dep.StoragePolicyClient() + dstPolicy, err := policyClient.GetPolicyByID(ctx, m.state.DstPolicyID) + if err != nil { + return task.StatusError, fmt.Errorf("failed to get destination policy: %w", err) + } + if dstPolicy == nil { + return task.StatusError, fmt.Errorf("destination policy %d does not exist (%w)", m.state.DstPolicyID, queue.CriticalErr) + } + if dstPolicy.Status == storagepolicy.StatusSuspended { + return task.StatusError, fmt.Errorf("destination policy %q is suspended (%w)", dstPolicy.Name, queue.CriticalErr) + } + + m.state.Phase = RelocatePhaseRelocating + if m.state.Failed == nil { + m.state.Failed = make(map[string]string) + } + + fm := manager.NewFileManager(dep, nil) + defer fm.Recycle() + + totalCount, totalSize := m.progressTotals(ctx, dep) + m.Lock() + m.progress[ProgressTypeRelocateCount] = &queue.Progress{Total: totalCount} + m.progress[ProgressTypeRelocateSize] = &queue.Progress{Total: totalSize} + m.Unlock() + + pending, err := m.pendingEntities(ctx, dep) + if err != nil { + return task.StatusError, fmt.Errorf("failed to list pending entities: %w", err) + } + + for len(pending) > 0 { + for _, e := range pending { + if err := ctx.Err(); err != nil { + m.ResumeAfter(0) + return task.StatusSuspending, nil + } + + if err := m.relocateEntity(ctx, dep, fm, dstPolicy, e); err != nil { + m.l.Warning("Failed to relocate entity %d: %s", e.ID, err) + m.state.Failed[strconv.Itoa(e.ID)] = err.Error() + } + m.state.Cursor = e.ID + } + + pending, err = m.pendingEntities(ctx, dep) + if err != nil { + return task.StatusError, fmt.Errorf("failed to list pending entities: %w", err) + } + } + + return task.StatusCompleted, nil +} + +// pendingEntities returns the next batch of entities still awaiting relocation. +// For explicit lists it walks the sorted IDs past the cursor; for policy +// migration it pages entities still assigned to the source policy. +func (m *RelocateTask) pendingEntities(ctx context.Context, dep dependency.Dep) ([]*ent.Entity, error) { + if len(m.state.EntityIDs) > 0 { + var ids []int + for _, id := range m.state.EntityIDs { + if id > m.state.Cursor { + ids = append(ids, id) + } + } + if len(ids) == 0 { + return nil, nil + } + if len(ids) > relocatePageSize { + ids = ids[:relocatePageSize] + } + return dep.DBClient().Entity.Query(). + Where(entity.IDIn(ids...)). + WithFile(). + WithUser(). + All(ctx) + } + + if m.state.SrcPolicyID == 0 { + return nil, nil + } + return dep.DBClient().Entity.Query(). + Where( + entity.StoragePolicyEntities(m.state.SrcPolicyID), + entity.IDGT(m.state.Cursor), + ). + Order(ent.Asc(entity.FieldID)). + Limit(relocatePageSize). + WithFile(). + WithUser(). + All(ctx) +} + +// progressTotals returns the expected entity count and byte total for the +// task's scope. Failures fall back to zero — progress is informational. +func (m *RelocateTask) progressTotals(ctx context.Context, dep dependency.Dep) (int64, int64) { + var q *ent.EntityQuery + if len(m.state.EntityIDs) > 0 { + q = dep.DBClient().Entity.Query().Where(entity.IDIn(m.state.EntityIDs...)) + } else if m.state.SrcPolicyID != 0 { + q = dep.DBClient().Entity.Query().Where(entity.StoragePolicyEntities(m.state.SrcPolicyID)) + } else { + return 0, 0 + } + + count, err := q.Clone().Count(ctx) + if err != nil { + return 0, 0 + } + var v []struct { + Sum int64 `json:"sum"` + } + if err := q.Clone().Aggregate(ent.Sum(entity.FieldSize)).Scan(ctx, &v); err != nil || len(v) == 0 { + return int64(count), 0 + } + return int64(count), v[0].Sum +} + +// relocateEntity streams one entity's blob to the destination policy, swaps its +// storage record, then deletes the old blob. A failure leaves the entity +// pointing at its original blob — the new object is dropped and recorded as a +// failure, so the entity remains readable. +func (m *RelocateTask) relocateEntity(ctx context.Context, dep dependency.Dep, fm manager.FileManager, dstPolicy *ent.StoragePolicy, e *ent.Entity) error { + if e.StoragePolicyEntities == m.state.DstPolicyID { + atomic.AddInt64(&m.progress[ProgressTypeRelocateCount].Current, 1) + return nil + } + + es, err := fm.GetEntitySource(ctx, e.ID) + if err != nil { + return fmt.Errorf("failed to open entity source: %w", err) + } + defer es.Close() + + // Generate the destination blob path from the policy's naming rules. + savePath := relocateSavePath(dstPolicy, e) + + // The source stream is already decrypted; wrap it again when the + // destination policy requires at-rest encryption. AES-CTR ciphertext + // matches the plaintext size, so Props.Size stays e.Size. + var encryptMetadata *types.EncryptMetadata + var reader io.ReadCloser = es + var seeker io.Seeker = es + if dstPolicy.Settings != nil && dstPolicy.Settings.Encryption { + cryptor, err := dep.EncryptorFactory(ctx)(types.CipherAES256CTR) + if err != nil { + return fmt.Errorf("failed to create cryptor: %w", err) + } + + encryptMetadata, err = cryptor.GenerateMetadata(ctx) + if err != nil { + return fmt.Errorf("failed to generate encrypt metadata: %w", err) + } + if err := cryptor.LoadMetadata(ctx, encryptMetadata); err != nil { + return fmt.Errorf("failed to load encrypt metadata: %w", err) + } + if err := cryptor.SetSource(es, es, e.Size, 0); err != nil { + return fmt.Errorf("failed to set cryptor source: %w", err) + } + reader = cryptor + seeker = cryptor + } + + req := &fs.UploadRequest{ + Props: &fs.UploadProps{ + Uri: relocateEntityUri(dep, e), + Size: e.Size, + SavePath: savePath, + }, + File: reader, + Seeker: seeker, + } + + dstDriver, err := fm.GetStorageDriver(ctx, fm.CastStoragePolicyOnSlave(ctx, dstPolicy)) + if err != nil { + return fmt.Errorf("failed to get destination driver: %w", err) + } + if err := dstDriver.Put(ctx, req); err != nil { + return fmt.Errorf("failed to write to destination policy: %w", err) + } + + // Swap the storage record; keep a copy of the old source path for cleanup. + oldSource := e.Source + oldPolicyID := e.StoragePolicyEntities + + newProps := types.EntityProps{} + if e.Props != nil { + newProps = *e.Props + } + newProps.EncryptMetadata = encryptMetadata + e.Props = &newProps + + if err := dep.FileClient().RelocateEntity(ctx, e, savePath, m.state.DstPolicyID); err != nil { + // The entity still points at its original blob — remove the orphaned + // copy on the destination so the entity stays consistent. + if _, derr := dstDriver.Delete(ctx, savePath); derr != nil { + m.l.Warning("Failed to clean up orphaned blob %q of entity %d: %s", savePath, e.ID, derr) + } + return fmt.Errorf("failed to update entity record: %w", err) + } + + // Remove the old blob; an orphan left behind is picked up by blob audit. + srcPolicy, err := dep.StoragePolicyClient().GetPolicyByID(ctx, oldPolicyID) + if err == nil && srcPolicy != nil { + if srcDriver, err := fm.GetStorageDriver(ctx, fm.CastStoragePolicyOnSlave(ctx, srcPolicy)); err == nil { + if failed, err := srcDriver.Delete(ctx, oldSource); err != nil || len(failed) > 0 { + m.l.Warning("Failed to delete old blob %q of entity %d: %s", oldSource, e.ID, err) + } + } + } + + atomic.AddInt64(&m.progress[ProgressTypeRelocateCount].Current, 1) + atomic.AddInt64(&m.progress[ProgressTypeRelocateSize].Current, e.Size) + return nil +} + +// relocateSavePath replicates the upload save-path rules (dir rule + name rule) +// for an entity being relocated. +func relocateSavePath(policy *ent.StoragePolicy, e *ent.Entity) string { + currentTime := time.Now() + originName := entityDisplayName(e) + dynamicReplace := func(rule string, pathAvailable bool) string { + return fs.ReplaceMagicVar(rule, fs.MagicVarProps{ + FsSeparator: fs.Separator, + PathAvailable: pathAvailable, + Time: currentTime, + UserID: entityOwnerID(e), + OriginName: originName, + }) + } + + dirRule := dynamicReplace(policy.DirNameRule, true) + nameRule := dynamicReplace(policy.FileNameRule, false) + + return path.Join(path.Clean(dirRule), nameRule) +} + +// relocateEntityUri builds a best-effort owner-scoped URI for the relocated +// blob's upload session record. The true file path is not reconstructed — the +// URI is informational on this path and never navigated. +func relocateEntityUri(dep dependency.Dep, e *ent.Entity) *fs.URI { + owner := hashid.EncodeUserID(dep.HashIDEncoder(), entityOwnerID(e)) + uri, err := fs.NewUriFromString(fs.NewMyUri(owner)) + if err != nil { + return nil + } + return uri.Join(entityDisplayName(e)) +} + +// entityDisplayName returns the display name of a file referencing the entity +// when the file edge is loaded, falling back to the blob's base name. +func entityDisplayName(e *ent.Entity) string { + if len(e.Edges.File) > 0 && e.Edges.File[0].Name != "" { + return e.Edges.File[0].Name + } + return path.Base(e.Source) +} + +func entityOwnerID(e *ent.Entity) int { + if e.Edges.User != nil { + return e.Edges.User.ID + } + return e.CreatedBy +} + +func (m *RelocateTask) Summarize(hasher hashid.Encoder) *queue.Summary { + if m.state == nil { + if err := json.Unmarshal([]byte(m.State()), &m.state); err != nil { + return nil + } + } + + props := map[string]any{ + SummaryKeyDstPolicy: m.state.DstPolicyID, + } + if m.state.SrcPolicyID != 0 { + props[SummaryKeySrcPolicy] = m.state.SrcPolicyID + } + if len(m.state.EntityIDs) > 0 { + props[SummaryKeyEntities] = len(m.state.EntityIDs) + } + if len(m.state.Failed) > 0 { + props[SummaryKeyRelocateFailed] = len(m.state.Failed) + } + + return &queue.Summary{ + NodeID: m.state.NodeID, + Phase: string(m.state.Phase), + Props: props, + } +} + +func (m *RelocateTask) Progress(ctx context.Context) queue.Progresses { + m.Lock() + defer m.Unlock() + return m.progress +} diff --git a/pkg/filemanager/workflows/relocate_test.go b/pkg/filemanager/workflows/relocate_test.go new file mode 100644 index 00000000..45f23022 --- /dev/null +++ b/pkg/filemanager/workflows/relocate_test.go @@ -0,0 +1,291 @@ +package workflows + +import ( + "context" + "encoding/json" + "fmt" + "os" + "path/filepath" + "testing" + + "github.com/cloudreve/Cloudreve/v4/application/constants" + "github.com/cloudreve/Cloudreve/v4/application/dependency" + "github.com/cloudreve/Cloudreve/v4/ent" + "github.com/cloudreve/Cloudreve/v4/ent/entity" + "github.com/cloudreve/Cloudreve/v4/ent/enttest" + "github.com/cloudreve/Cloudreve/v4/ent/storagepolicy" + "github.com/cloudreve/Cloudreve/v4/ent/task" + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/boolset" + "github.com/cloudreve/Cloudreve/v4/pkg/cache" + "github.com/cloudreve/Cloudreve/v4/pkg/conf" + "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" + "github.com/cloudreve/Cloudreve/v4/pkg/util" + "github.com/stretchr/testify/require" +) + +func newRelocateTestDep(t *testing.T, client *ent.Client) (context.Context, dependency.Dep) { + logger := logging.NewConsoleLogger(logging.LevelError) + cfg, err := conf.NewIniConfigProvider(t.TempDir()+"/conf.ini", logger) + require.NoError(t, err) + hasher, err := hashid.New("test-salt") + require.NoError(t, err) + + dep := dependency.NewDependency( + dependency.WithDbClient(client), + dependency.WithConfigProvider(cfg), + dependency.WithKV(cache.NewMemoStore("", logger)), + dependency.WithLogger(logger), + dependency.WithHashIDEncoder(hasher), + ) + ctx := context.WithValue(context.Background(), dependency.DepCtx{}, dep) + return ctx, dep +} + +func relocateTaskWithState(t *testing.T, state *RelocateTaskState) *RelocateTask { + stateBytes, err := json.Marshal(state) + require.NoError(t, err) + return &RelocateTask{ + DBTask: &queue.DBTask{ + Task: &ent.Task{ + Type: queue.RelocateTaskType, + PrivateState: string(stateBytes), + PublicState: &types.TaskPublicState{}, + }, + }, + } +} + +func TestNewRelocateTaskSortsIDs(t *testing.T) { + tk, err := NewRelocateTask(context.Background(), []int{9, 1, 5}, 3) + require.NoError(t, err) + + rt := tk.(*RelocateTask) + state := &RelocateTaskState{} + require.NoError(t, json.Unmarshal([]byte(rt.Task.PrivateState), state)) + require.Equal(t, []int{1, 5, 9}, state.EntityIDs) + require.Equal(t, 3, state.DstPolicyID) + require.Equal(t, queue.RelocateTaskType, rt.Task.Type) +} + +func TestRelocateRejectsIdenticalPolicies(t *testing.T) { + m := relocateTaskWithState(t, &RelocateTaskState{SrcPolicyID: 4, DstPolicyID: 4}) + m.state = &RelocateTaskState{SrcPolicyID: 4, DstPolicyID: 4} + + status, err := m.relocate(context.Background(), nil) + require.Error(t, err) + require.Equal(t, task.StatusError, status) +} + +func TestRelocatePendingExplicitList(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + ctx, dep := newRelocateTestDep(t, client) + + p := client.StoragePolicy.Create().SetName("src").SetType("local"). + SetStatus(storagepolicy.StatusActive).SaveX(ctx) + mk := func(src string, size int64) int { + return client.Entity.Create().SetType(0).SetSource(src).SetSize(size). + SetStoragePolicyEntities(p.ID).SaveX(ctx).ID + } + a, b, c := mk("a", 1), mk("b", 2), mk("c", 3) + + m := relocateTaskWithState(t, &RelocateTaskState{ + EntityIDs: []int{c, a, b}, + DstPolicyID: p.ID + 1, + }) + m.state = &RelocateTaskState{EntityIDs: []int{a, b, c}, DstPolicyID: p.ID + 1} + + pending, err := m.pendingEntities(ctx, dep) + require.NoError(t, err) + require.Len(t, pending, 3) + + // Cursor resumes after the last processed ID. + m.state.Cursor = b + pending, err = m.pendingEntities(ctx, dep) + require.NoError(t, err) + require.Len(t, pending, 1) + require.Equal(t, c, pending[0].ID) +} + +func TestRelocatePendingPolicyScope(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + ctx, dep := newRelocateTestDep(t, client) + + src := client.StoragePolicy.Create().SetName("src").SetType("local"). + SetStatus(storagepolicy.StatusActive).SaveX(ctx) + dst := client.StoragePolicy.Create().SetName("dst").SetType("local"). + SetStatus(storagepolicy.StatusActive).SaveX(ctx) + + stay := client.Entity.Create().SetType(0).SetSource("s").SetSize(1). + SetStoragePolicyEntities(src.ID).SaveX(ctx) + moved := client.Entity.Create().SetType(0).SetSource("m").SetSize(1). + SetStoragePolicyEntities(dst.ID).SaveX(ctx) + + m := relocateTaskWithState(t, &RelocateTaskState{SrcPolicyID: src.ID, DstPolicyID: dst.ID}) + m.state = &RelocateTaskState{SrcPolicyID: src.ID, DstPolicyID: dst.ID} + + pending, err := m.pendingEntities(ctx, dep) + require.NoError(t, err) + require.Len(t, pending, 1) + require.Equal(t, stay.ID, pending[0].ID) + + // After relocation the entity leaves the source-policy scope. + m.state.Cursor = stay.ID + pending, err = m.pendingEntities(ctx, dep) + require.NoError(t, err) + require.Empty(t, pending) + _ = moved +} + +func TestRelocateProgressTotals(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + ctx, dep := newRelocateTestDep(t, client) + + src := client.StoragePolicy.Create().SetName("src").SetType("local"). + SetStatus(storagepolicy.StatusActive).SaveX(ctx) + other := client.StoragePolicy.Create().SetName("other").SetType("local"). + SetStatus(storagepolicy.StatusActive).SaveX(ctx) + + mk := func(policyID int, size int64) int { + return client.Entity.Create().SetType(0).SetSource("s").SetSize(size). + SetStoragePolicyEntities(policyID).SaveX(ctx).ID + } + a := mk(src.ID, 10) + b := mk(src.ID, 20) + mk(other.ID, 99) + + m := relocateTaskWithState(t, &RelocateTaskState{SrcPolicyID: src.ID, DstPolicyID: other.ID}) + m.state = &RelocateTaskState{SrcPolicyID: src.ID, DstPolicyID: other.ID} + count, size := m.progressTotals(ctx, dep) + require.Equal(t, int64(2), count) + require.Equal(t, int64(30), size) + + m.state = &RelocateTaskState{EntityIDs: []int{a, b}, DstPolicyID: other.ID} + count, size = m.progressTotals(ctx, dep) + require.Equal(t, int64(2), count) + require.Equal(t, int64(30), size) +} + +func TestRelocateSavePathNamingRules(t *testing.T) { + p := &ent.StoragePolicy{ + DirNameRule: "vault/{uid}", + FileNameRule: "{originname}", + } + e := &ent.Entity{Source: "uploads/blob.bin"} + e.Edges.File = []*ent.File{{Name: "report.pdf"}} + e.Edges.User = &ent.User{ID: 7} + + require.Equal(t, "vault/7/report.pdf", relocateSavePath(p, e)) +} + +func TestRelocateEntityUriOwnerScoped(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + ctx, dep := newRelocateTestDep(t, client) + _ = ctx + + e := &ent.Entity{Source: "uploads/blob.bin"} + e.Edges.File = []*ent.File{{Name: "report.pdf"}} + e.Edges.User = &ent.User{ID: 7} + + uri := relocateEntityUri(dep, e) + require.NotNil(t, uri) + require.Equal(t, "report.pdf", uri.Name()) + require.Equal(t, constants.FileSystemMy, uri.FileSystem()) + require.NotEmpty(t, uri.ID("")) +} + +func TestRelocateSummarize(t *testing.T) { + m := relocateTaskWithState(t, &RelocateTaskState{ + SrcPolicyID: 1, + DstPolicyID: 2, + Failed: map[string]string{"5": "io error"}, + }) + s := m.Summarize(nil) + require.NotNil(t, s) + require.Equal(t, 2, s.Props[SummaryKeyDstPolicy]) + require.Equal(t, 1, s.Props[SummaryKeySrcPolicy]) + require.Equal(t, 1, s.Props[SummaryKeyRelocateFailed]) +} + +func TestRelocateEntitySkipsSamePolicy(t *testing.T) { + m := relocateTaskWithState(t, &RelocateTaskState{DstPolicyID: 3}) + m.state = &RelocateTaskState{DstPolicyID: 3} + m.progress = queue.Progresses{ProgressTypeRelocateCount: &queue.Progress{}} + + e := &ent.Entity{ID: 1, StoragePolicyEntities: 3} + require.NoError(t, m.relocateEntity(context.Background(), nil, nil, nil, e)) + require.Equal(t, int64(1), m.progress[ProgressTypeRelocateCount].Current) +} + +// TestRelocateEntityLocalEndToEnd streams a real blob between two local +// policies: content lands under the destination naming rules, the entity +// record swaps to the new policy, and the old blob is removed. +func TestRelocateEntityLocalEndToEnd(t *testing.T) { + tmp := t.TempDir() + oldWd, err := os.Getwd() + require.NoError(t, err) + require.NoError(t, os.Chdir(tmp)) + oldUseWd := util.UseWorkingDir + util.UseWorkingDir = true + t.Cleanup(func() { + util.UseWorkingDir = oldUseWd + require.NoError(t, os.Chdir(oldWd)) + }) + + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + ctx, dep := newRelocateTestDep(t, client) + + group := client.Group.Create().SetName("g").SetPermissions(&boolset.BooleanSet{}).SaveX(ctx) + user := client.User.Create().SetEmail("u@example.com").SetNick("u").SetGroup(group).SaveX(ctx) + + srcPolicy := client.StoragePolicy.Create().SetName("src").SetType("local"). + SetStatus(storagepolicy.StatusActive).SaveX(ctx) + dstPolicy := client.StoragePolicy.Create().SetName("dst").SetType("local"). + SetStatus(storagepolicy.StatusActive). + SetDirNameRule("vault/{uid}").SetFileNameRule("{originname}").SaveX(ctx) + + content := []byte("relocate-me") + srcFile := filepath.Join(tmp, "blob.bin") + require.NoError(t, os.WriteFile(srcFile, content, 0o644)) + + e := client.Entity.Create(). + SetType(0). + SetSource(srcFile). + SetSize(int64(len(content))). + SetReferenceCount(1). + SetStoragePolicyEntities(srcPolicy.ID). + SetCreatedBy(user.ID). + SaveX(ctx) + e = client.Entity.Query().Where(entity.ID(e.ID)).WithFile().WithUser().OnlyX(ctx) + + fm := manager.NewFileManager(dep, user) + defer fm.Recycle() + + m := relocateTaskWithState(t, &RelocateTaskState{DstPolicyID: dstPolicy.ID}) + m.state = &RelocateTaskState{DstPolicyID: dstPolicy.ID} + m.l = logging.NewConsoleLogger(logging.LevelError) + m.progress = queue.Progresses{ + ProgressTypeRelocateCount: &queue.Progress{}, + ProgressTypeRelocateSize: &queue.Progress{}, + } + + require.NoError(t, m.relocateEntity(ctx, dep, fm, dstPolicy, e)) + + updated := client.Entity.GetX(ctx, e.ID) + require.Equal(t, dstPolicy.ID, updated.StoragePolicyEntities) + require.Equal(t, fmt.Sprintf("vault/%d/blob.bin", user.ID), updated.Source) + require.NoFileExists(t, srcFile) + got, err := os.ReadFile(filepath.Join(tmp, updated.Source)) + require.NoError(t, err) + require.Equal(t, content, got) + require.Equal(t, int64(1), m.progress[ProgressTypeRelocateCount].Current) + require.Equal(t, int64(len(content)), m.progress[ProgressTypeRelocateSize].Current) +} diff --git a/routers/controllers/admin.go b/routers/controllers/admin.go index 852a2c83..b6d2e4a8 100644 --- a/routers/controllers/admin.go +++ b/routers/controllers/admin.go @@ -485,6 +485,15 @@ func AdminBatchDeleteEntity(c *gin.Context) { } } +func AdminRelocateEntity(c *gin.Context) { + service := ParametersFromContext[*admin.RelocateEntityService](c, admin.RelocateEntityParamCtx{}) + res, err := service.Relocate(c) + if respondErr(c, err) { + return + } + c.JSON(200, serializer.Response{Data: res}) +} + func AdminCleanupTask(c *gin.Context) { service := ParametersFromContext[*admin.CleanupTaskService](c, admin.CleanupTaskParameterCtx{}) err := service.CleanupTask(c) diff --git a/routers/router.go b/routers/router.go index f573268b..001b4225 100644 --- a/routers/router.go +++ b/routers/router.go @@ -1276,6 +1276,12 @@ func initMasterRouter(dep dependency.Dep) *gin.Engine { controllers.FromJSON[adminsvc.BatchEntityService](adminsvc.BatchEntityParamCtx{}), controllers.AdminBatchDeleteEntity, ) + // Relocate entities to another storage policy + entity.POST("relocate", + middleware.RequiredScopes(types.ScopeAdminWrite), + controllers.FromJSON[adminsvc.RelocateEntityService](adminsvc.RelocateEntityParamCtx{}), + controllers.AdminRelocateEntity, + ) // Get entity url entity.GET("url/:id", controllers.FromUri[adminsvc.SingleEntityService](adminsvc.SingleEntityParamCtx{}), diff --git a/service/admin/file.go b/service/admin/file.go index 4fd7ecea..4c936800 100644 --- a/service/admin/file.go +++ b/service/admin/file.go @@ -14,7 +14,9 @@ import ( "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/manager" "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/manager/entitysource" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/workflows" "github.com/cloudreve/Cloudreve/v4/pkg/hashid" + "github.com/cloudreve/Cloudreve/v4/pkg/queue" "github.com/cloudreve/Cloudreve/v4/pkg/serializer" "github.com/cloudreve/Cloudreve/v4/pkg/setting" "github.com/gin-gonic/gin" @@ -538,6 +540,62 @@ func (s *BatchEntityService) Delete(c *gin.Context) error { return nil } +type ( + // RelocateEntityService moves blobs between storage policies. Exactly one + // scope is required: explicit entity IDs, or a source policy whose entities + // are all migrated. + RelocateEntityService struct { + EntityIDs []int `json:"entity_ids"` + SrcPolicyID int `json:"src_policy_id"` + DstPolicyID int `json:"dst_policy_id" binding:"required"` + } + RelocateEntityParamCtx struct{} +) + +func (s *RelocateEntityService) Relocate(c *gin.Context) (*RelocateTaskResponse, error) { + dep := dependency.FromContext(c) + hasher := dep.HashIDEncoder() + + if len(s.EntityIDs) == 0 && s.SrcPolicyID == 0 { + return nil, serializer.NewError(serializer.CodeParamErr, "either entity_ids or src_policy_id is required", nil) + } + if len(s.EntityIDs) > 0 && s.SrcPolicyID != 0 { + return nil, serializer.NewError(serializer.CodeParamErr, "entity_ids and src_policy_id are mutually exclusive", nil) + } + if s.SrcPolicyID == s.DstPolicyID && s.SrcPolicyID != 0 { + return nil, serializer.NewError(serializer.CodeParamErr, "source and destination policies are identical", nil) + } + + dstPolicy, err := dep.StoragePolicyClient().GetPolicyByID(c, s.DstPolicyID) + if err != nil || dstPolicy == nil { + return nil, serializer.NewError(serializer.CodeParamErr, "destination policy does not exist", err) + } + + var t queue.Task + if s.SrcPolicyID != 0 { + srcPolicy, err := dep.StoragePolicyClient().GetPolicyByID(c, s.SrcPolicyID) + if err != nil || srcPolicy == nil { + return nil, serializer.NewError(serializer.CodeParamErr, "source policy does not exist", err) + } + t, err = workflows.NewRelocatePolicyTask(c, s.SrcPolicyID, s.DstPolicyID) + } else { + t, err = workflows.NewRelocateTask(c, s.EntityIDs, s.DstPolicyID) + } + 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 &RelocateTaskResponse{ID: hashid.EncodeTaskID(hasher, t.ID())}, nil +} + +type RelocateTaskResponse struct { + ID string `json:"id"` +} + func (s *SingleEntityService) Url(c *gin.Context) (string, error) { dep := dependency.FromContext(c) fileClient := dep.FileClient()