feat(admin): storage-policy relocation workflow (#9 #125 #136)pull/3587/head
commit
0aed15de29
@ -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