Resumable RelocateTask moves entity blobs between storage policies: streams plaintext through EntitySource, re-encrypts when the destination policy requires it, swaps source/policy/props only after a successful Put, deletes the old blob after the metadata swap, and cleans up the orphaned destination blob when the swap fails. Supports an explicit entity list or whole-source-policy migration with cursor-based resume and per-entity failure tracking. Admin endpoint POST /admin/entity/relocate validates the scope (entity_ids XOR src_policy_id), rejects missing/suspended/identical destination policies, and queues the task on IoIntenseQueue. Admin UI: batch Relocate action on the entity list and a per-policy migrate action on storage policy cards, backed by a shared RelocateDialog + SinglePolicySelectionInput. Resolves upstream cloudreve/cloudreve#2262 (#9), #3518 (#125), and #3576 (#136). Authored By: TDvorak <info@tdvorak.dev> Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>pull/3582/head
parent
73b956e4b3
commit
a8c06154fc
@ -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<number | undefined>(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 (
|
||||||
|
<Dialog open={open} onClose={onClose} maxWidth="xs" fullWidth>
|
||||||
|
<DialogTitle>{t("entity.relocateTitle")}</DialogTitle>
|
||||||
|
<DialogContent>
|
||||||
|
<DialogContentText sx={{ mb: 2 }}>
|
||||||
|
{srcPolicyID
|
||||||
|
? t("entity.relocatePolicyDes", { name: srcPolicyName ?? "" })
|
||||||
|
: t("entity.relocateEntitiesDes", { num: entityIDs?.length ?? 0 })}
|
||||||
|
</DialogContentText>
|
||||||
|
<SinglePolicySelectionInput value={dstPolicyID} onChange={setDstPolicyID} />
|
||||||
|
</DialogContent>
|
||||||
|
<DialogActions>
|
||||||
|
<Button onClick={onClose}>{t("common:cancel")}</Button>
|
||||||
|
<Button variant="contained" onClick={onSubmit} disabled={loading || !dstPolicyID}>
|
||||||
|
{t("entity.relocateStart")}
|
||||||
|
</Button>
|
||||||
|
</DialogActions>
|
||||||
|
</Dialog>
|
||||||
|
);
|
||||||
|
};
|
||||||
|
|
||||||
|
export default RelocateDialog;
|
||||||
@ -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
|
||||||
|
}
|
||||||
@ -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)
|
||||||
|
}
|
||||||
Loading…
Reference in new issue