diff --git a/pkg/filemanager/workflows/archive.go b/pkg/filemanager/workflows/archive.go index 53fe21b1..2bef70b5 100644 --- a/pkg/filemanager/workflows/archive.go +++ b/pkg/filemanager/workflows/archive.go @@ -125,6 +125,11 @@ func (m *CreateArchiveTask) Do(ctx context.Context) (task.Status, error) { } m.state = state + // Prefer the node hosting the first source file's storage policy. + if len(m.state.Uris) > 0 { + preferPolicyNode(ctx, dep, &m.state.NodeState, m.state.Uris[0]) + } + // select node node, err := allocateNode(ctx, dep, &m.state.NodeState, types.NodeCapabilityCreateArchive) if err != nil { diff --git a/pkg/filemanager/workflows/extract.go b/pkg/filemanager/workflows/extract.go index 3c379cdf..b145a0aa 100644 --- a/pkg/filemanager/workflows/extract.go +++ b/pkg/filemanager/workflows/extract.go @@ -126,6 +126,9 @@ func (m *ExtractArchiveTask) Do(ctx context.Context) (task.Status, error) { } m.state = state + // Prefer the node hosting the source file's storage policy. + preferPolicyNode(ctx, dep, &m.state.NodeState, m.state.Uri) + // select node node, err := allocateNode(ctx, dep, &m.state.NodeState, types.NodeCapabilityExtractArchive) if err != nil { diff --git a/pkg/filemanager/workflows/worfklows.go b/pkg/filemanager/workflows/worfklows.go index 475dd9be..1356d099 100644 --- a/pkg/filemanager/workflows/worfklows.go +++ b/pkg/filemanager/workflows/worfklows.go @@ -9,7 +9,11 @@ import ( "github.com/cloudreve/Cloudreve/v4/application/dependency" "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/inventory" "github.com/cloudreve/Cloudreve/v4/pkg/cluster" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs/dbfs" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/manager" "github.com/cloudreve/Cloudreve/v4/pkg/queue" "github.com/cloudreve/Cloudreve/v4/pkg/util" ) @@ -60,3 +64,38 @@ func prepareTempFolder(ctx context.Context, dep dependency.Dep, t queue.Task) (s dep.Logger().Info("Temp folder created: %s", tempPath) return tempPath, nil } + +// preferPolicyNode seeds the preferred node with the node hosting the given +// file's storage policy, so data-heavy tasks run where the blob already +// lives. No-op when the URI is invalid, has no entity, or the policy is not +// bound to a node — the pool then falls back to weighted selection. +func preferPolicyNode(ctx context.Context, dep dependency.Dep, state *NodeState, uriStr string) { + if state.NodeID > 0 { + return + } + + uri, err := fs.NewUriFromString(uriStr) + if err != nil { + return + } + + fm := manager.NewFileManager(dep, inventory.UserFromContext(ctx)) + defer fm.Recycle() + + file, err := fm.Get(ctx, uri, dbfs.WithFileEntities()) + if err != nil || file == nil { + return + } + + entity := file.PrimaryEntity() + if entity == nil { + return + } + + policy, err := dep.StoragePolicyClient().GetPolicyByID(ctx, entity.PolicyID()) + if err != nil || policy == nil { + return + } + + state.NodeID = policy.NodeID +}