From 8febe344ac001ed6a6644ffdb580a17cbdc4c1fa Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Tom=C3=A1=C5=A1=20Dvo=C5=99=C3=A1k?=
<150935816+Dvorinka@users.noreply.github.com>
Date: Sun, 20 Sep 2026 00:36:07 +0200
Subject: [PATCH] feat: group-level task node pools and user node selection
(#186)
Groups can restrict task dispatch to an allowed node pool and optionally
let users pick a target node for remote downloads and archive ops.
- GroupSetting: allowed_nodes (empty = all) + allow_select_node toggle
- NodePool.Get filters candidates by the allowed pool; explicit picks
outside the pool fall back to weighted selection within it
- NodeState persists the pool so queued tasks keep their restrictions
- resolveNodeSelection validates picks: group toggle, pool membership,
active status, task capability; node hashids on the wire
- Site config exposes task_nodes + allow_select_node for permitted groups
- Admin group editor: functional node multi-select + allow-select switch
- Task dialogs: shared TargetNodeSelect on remote download, archive
create and extract
Generated with [Devin](https://devin.ai)
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
---
ROADMAP.md | 5 +-
frontend/src/api/dashboard.ts | 2 +
frontend/src/api/site.ts | 2 +
frontend/src/api/workflow.ts | 2 +
.../Group/EditGroup/FileManagementSection.tsx | 22 ++-
.../EditGroup/MultipleNodeSelectionInput.tsx | 110 +++++++++++----
.../Common/Form/TargetNodeSelect.tsx | 48 +++++++
.../FileManager/Dialogs/CreateArchive.tsx | 7 +-
.../Dialogs/CreateRemoteDownload.tsx | 7 +-
.../FileManager/Dialogs/ExtractArchive.tsx | 16 ++-
inventory/types/types.go | 6 +
middleware/cluster.go | 2 +-
pkg/cluster/pool.go | 18 ++-
pkg/cluster/pool_test.go | 44 +++++-
pkg/filemanager/workflows/archive.go | 4 +-
pkg/filemanager/workflows/extract.go | 28 ++--
pkg/filemanager/workflows/remote_download.go | 3 +
pkg/filemanager/workflows/upload.go | 2 +-
pkg/filemanager/workflows/workflows.go | 17 ++-
service/basic/site.go | 56 ++++++++
service/explorer/workflows.go | 84 +++++++++---
service/explorer/workflows_nodeselect_test.go | 127 ++++++++++++++++++
22 files changed, 535 insertions(+), 77 deletions(-)
create mode 100644 frontend/src/component/Common/Form/TargetNodeSelect.tsx
create mode 100644 service/explorer/workflows_nodeselect_test.go
diff --git a/ROADMAP.md b/ROADMAP.md
index db9ef1ff..0eb9285f 100644
--- a/ROADMAP.md
+++ b/ROADMAP.md
@@ -191,9 +191,10 @@ Order = user-visible value first; each ships with backend + UI + tests.
- [x] PR #144 — task `creator_ip` capture with CIDR-capable admin filter (#115 OSS half), group remote-download quotas per count + per volume (#16), yt-dlp downloader provider (#88), progressive image preview (#113), v3 migrator `DatabaseURL` passthrough (#42)
- [x] `activity_event` entity (immutable, tx-aware, actor+IP+CID) + per-file Activity dialog + admin `/admin/event` feed + per-type enablement + retention cron (#184)
- [x] Coverage wave 2: email/user-activated/token-refresh/share-viewed/version/metadata/view/thumb/live-photo/copy-from/webdav/profile+security/oauth/admin-ops/import (1bbaddf)
- - [ ] Event coverage remainder (needs unbuilt features): payment_*, link/unlink_account, membership_unsubscribe, report_abuse, mount, quota-notify
+ - [ ] Event coverage remainder (needs unbuilt features): payment_*, link/unlink_account, membership_unsubscribe, mount, quota-notify
- [x] site announcement: `announcement` setting (markdown) + post-login modal + per-user dismissal re-triggering on content change (#184)
- - [ ] group `allowed_nodes` + task `target_node`; `abuse_report` + admin queue + share context-menu Report entry
+ - [x] `abuse_report` entity + public `POST /abuse/report` (IP rate-limit + `abuse_captcha` gate) + admin `/admin/abuse` queue (resolve/dismiss + reversible share block) + share-menu Report entry (#185)
+ - [x] group `allowed_nodes` pool + `allow_select_node` + task `target_node` dispatch (persisted in task state, weighted LB within pool); group admin multi-select + task-dialog node picker
## 5. Phase C — security + quality
diff --git a/frontend/src/api/dashboard.ts b/frontend/src/api/dashboard.ts
index 962c1e92..60e13bd0 100644
--- a/frontend/src/api/dashboard.ts
+++ b/frontend/src/api/dashboard.ts
@@ -66,6 +66,8 @@ export interface GroupSetting {
redirected_source?: boolean;
login_ip_whitelist?: string[];
default_pinned?: number[];
+ allowed_nodes?: number[];
+ allow_select_node?: boolean;
}
export interface AdminListGroupResponse {
diff --git a/frontend/src/api/site.ts b/frontend/src/api/site.ts
index caf9e8a2..e227484d 100644
--- a/frontend/src/api/site.ts
+++ b/frontend/src/api/site.ts
@@ -34,6 +34,8 @@ export interface SiteConfig {
sso_auto_redirect?: boolean;
download_cdn_routes?: { name: string; url: string }[];
abuse_captcha?: boolean;
+ allow_select_node?: boolean;
+ task_nodes?: { id: string; name: string }[];
logo?: string;
logo_light?: string;
tos_url?: string;
diff --git a/frontend/src/api/workflow.ts b/frontend/src/api/workflow.ts
index e28ec1aa..e0b7004e 100644
--- a/frontend/src/api/workflow.ts
+++ b/frontend/src/api/workflow.ts
@@ -6,6 +6,7 @@ export interface ArchiveWorkflowService {
encoding?: string;
password?: string;
file_mask?: string[];
+ target_node?: string;
}
export interface TaskListResponse {
@@ -106,6 +107,7 @@ export interface DownloadWorkflowService {
password?: string;
headers?: string[];
provider?: string;
+ target_node?: string;
}
export interface ImportWorkflowService {
diff --git a/frontend/src/component/Admin/Group/EditGroup/FileManagementSection.tsx b/frontend/src/component/Admin/Group/EditGroup/FileManagementSection.tsx
index 102c5488..88457ce9 100644
--- a/frontend/src/component/Admin/Group/EditGroup/FileManagementSection.tsx
+++ b/frontend/src/component/Admin/Group/EditGroup/FileManagementSection.tsx
@@ -391,14 +391,32 @@ const FileManagementSection = () => {
-
+
+ setGroup((p: GroupEnt) => ({
+ ...p,
+ settings: { ...p.settings, allowed_nodes: v.length > 0 ? v : undefined },
+ }))
+ }
+ />
{t("group.allowedNodesDes")}
}
+ control={
+
+ setGroup((p: GroupEnt) => ({
+ ...p,
+ settings: { ...p.settings, allow_select_node: e.target.checked ? true : undefined },
+ }))
+ }
+ />
+ }
label={
{t("group.allowSelectNode")}
diff --git a/frontend/src/component/Admin/Group/EditGroup/MultipleNodeSelectionInput.tsx b/frontend/src/component/Admin/Group/EditGroup/MultipleNodeSelectionInput.tsx
index 2be4da18..9240da0b 100644
--- a/frontend/src/component/Admin/Group/EditGroup/MultipleNodeSelectionInput.tsx
+++ b/frontend/src/component/Admin/Group/EditGroup/MultipleNodeSelectionInput.tsx
@@ -1,40 +1,92 @@
-import { ListItemText } from "@mui/material";
+import { Box, Checkbox, FormControl, ListItemText, SelectChangeEvent } from "@mui/material";
+import { useEffect, useState } from "react";
import { useTranslation } from "react-i18next";
+import { getNodeList } from "../../../../api/api";
+import { Node } from "../../../../api/dashboard";
import { useAppDispatch } from "../../../../redux/hooks";
-import { DenseSelect } from "../../../Common/StyledComponents";
+import { DenseSelect, SquareChip } from "../../../Common/StyledComponents";
+import { SquareMenuItem } from "../../../FileManager/ContextMenu/ContextMenu";
-const MultipleNodeSelectionInput = () => {
+export interface MultipleNodeSelectionInputProps {
+ value: number[];
+ onChange: (value: number[]) => void;
+}
+
+// MultipleNodeSelectionInput picks the pool of nodes a group's tasks may run
+// on. Empty means all nodes are eligible.
+const MultipleNodeSelectionInput = ({ value, onChange }: MultipleNodeSelectionInputProps) => {
const { t } = useTranslation("dashboard");
const dispatch = useAppDispatch();
+ const [nodes, setNodes] = useState([]);
+ const [loading, setLoading] = useState(false);
+ const [nodeMap, setNodeMap] = useState>({});
- return (
- ) => {
+ const {
+ target: { value: v },
+ } = event;
+ onChange(typeof v === "string" ? v.split(",").map((x) => parseInt(x)) : (v as number[]));
+ };
+
+ useEffect(() => {
+ setLoading(true);
+ dispatch(getNodeList({ page: 1, page_size: 1000, order_by: "id", order_direction: "asc" }))
+ .then((res) => {
+ setNodes(res.nodes);
+ setNodeMap(
+ res.nodes.reduce(
+ (acc, n) => {
+ acc[n.id] = n;
+ return acc;
},
- },
- },
- }}
- renderValue={(selected) => {
- return (
- {t("group.allNodes")}}
- slotProps={{
- primary: { color: "textSecondary", variant: "body2" },
- }}
- />
+ {} as Record,
+ ),
);
- }}
- >
+ })
+ .finally(() => {
+ setLoading(false);
+ });
+ }, []);
+
+ return (
+
+
+ (selected as number[]).length === 0 ? (
+ {t("group.allNodes")}}
+ slotProps={{
+ primary: { color: "textSecondary", variant: "body2" },
+ }}
+ />
+ ) : (
+
+ {(selected as number[]).map((id) => (
+
+ ))}
+
+ )
+ }
+ >
+ {nodes.map((n) => (
+
+ -1} />
+
+
+ ))}
+
+
);
};
diff --git a/frontend/src/component/Common/Form/TargetNodeSelect.tsx b/frontend/src/component/Common/Form/TargetNodeSelect.tsx
new file mode 100644
index 00000000..5f8895a0
--- /dev/null
+++ b/frontend/src/component/Common/Form/TargetNodeSelect.tsx
@@ -0,0 +1,48 @@
+import { FormControl, InputLabel, MenuItem, Select } from "@mui/material";
+import { useTranslation } from "react-i18next";
+import { useAppSelector } from "../../../redux/hooks.ts";
+
+export interface TargetNodeSelectProps {
+ value: string;
+ onChange: (value: string) => void;
+}
+
+export const useShowTargetNodeSelect = () =>
+ useAppSelector(
+ (state) =>
+ !!state.siteConfig.explorer?.config?.allow_select_node &&
+ (state.siteConfig.explorer?.config?.task_nodes?.length ?? 0) > 0,
+ );
+
+const TargetNodeSelect = ({ value, onChange }: TargetNodeSelectProps) => {
+ const { t } = useTranslation();
+ const show = useShowTargetNodeSelect();
+ const taskNodes = useAppSelector((state) => state.siteConfig.explorer?.config?.task_nodes);
+
+ if (!show) {
+ return null;
+ }
+
+ return (
+
+ {t("application:modals.processNode")}
+
+
+ );
+};
+
+export default TargetNodeSelect;
diff --git a/frontend/src/component/FileManager/Dialogs/CreateArchive.tsx b/frontend/src/component/FileManager/Dialogs/CreateArchive.tsx
index 6227ea0f..3fb5b83b 100644
--- a/frontend/src/component/FileManager/Dialogs/CreateArchive.tsx
+++ b/frontend/src/component/FileManager/Dialogs/CreateArchive.tsx
@@ -9,6 +9,7 @@ import { getFileLinkedUri } from "../../../util";
import CrUri from "../../../util/uri.ts";
import { OutlineIconTextField } from "../../Common/Form/OutlineIconTextField.tsx";
import { PathSelectorForm } from "../../Common/Form/PathSelectorForm.tsx";
+import TargetNodeSelect from "../../Common/Form/TargetNodeSelect.tsx";
import { ViewTaskAction } from "../../Common/Snackbar/snackbar.tsx";
import DraggableDialog from "../../Dialogs/DraggableDialog.tsx";
import Archive from "../../Icons/Archive.tsx";
@@ -24,6 +25,7 @@ const CreateArchive = () => {
const [loading, setLoading] = useState(false);
const [fileName, setFileName] = useState("archive.zip");
const [path, setPath] = useState("");
+ const [targetNode, setTargetNode] = useState("");
const open = useAppSelector((state) => state.globalState.createArchiveDialogOpen);
const targets = useAppSelector((state) => state.globalState.createArchiveDialogFiles);
@@ -32,6 +34,7 @@ const CreateArchive = () => {
useEffect(() => {
if (open) {
setPath(current ?? "");
+ setTargetNode("");
}
}, [open]);
@@ -50,6 +53,7 @@ const CreateArchive = () => {
sendCreateArchive({
src: targets?.map((t) => getFileLinkedUri(t)),
dst: dst.join(fileName).toString(),
+ target_node: targetNode || undefined,
}),
)
.then(() => {
@@ -63,7 +67,7 @@ const CreateArchive = () => {
.finally(() => {
setLoading(false);
});
- }, [targets, fileName, path]);
+ }, [targets, fileName, path, targetNode]);
return (
{
+
diff --git a/frontend/src/component/FileManager/Dialogs/CreateRemoteDownload.tsx b/frontend/src/component/FileManager/Dialogs/CreateRemoteDownload.tsx
index 00845748..3f592f51 100644
--- a/frontend/src/component/FileManager/Dialogs/CreateRemoteDownload.tsx
+++ b/frontend/src/component/FileManager/Dialogs/CreateRemoteDownload.tsx
@@ -11,6 +11,7 @@ import CrUri, { Filesystem } from "../../../util/uri.ts";
import { FileDisplayForm } from "../../Common/Form/FileDisplayForm.tsx";
import { OutlineIconTextField } from "../../Common/Form/OutlineIconTextField.tsx";
import { PathSelectorForm } from "../../Common/Form/PathSelectorForm.tsx";
+import TargetNodeSelect from "../../Common/Form/TargetNodeSelect.tsx";
import { ViewTaskAction } from "../../Common/Snackbar/snackbar.tsx";
import DraggableDialog from "../../Dialogs/DraggableDialog.tsx";
import Edit from "../../Icons/Edit.tsx";
@@ -41,6 +42,7 @@ const CreateRemoteDownload = () => {
const [password, setPassword] = useState("");
const [headers, setHeaders] = useState("");
const [provider, setProvider] = useState("");
+ const [targetNode, setTargetNode] = useState("");
const providers = useAppSelector((state) => state.siteConfig.explorer?.config?.remote_download_providers);
const open = useAppSelector((state) => state.globalState.remoteDownloadDialogOpen);
@@ -58,6 +60,7 @@ const CreateRemoteDownload = () => {
setPassword("");
setHeaders("");
setProvider("");
+ setTargetNode("");
}
}, [open]);
@@ -81,6 +84,7 @@ const CreateRemoteDownload = () => {
password: password || undefined,
headers: headers ? headers.split("\n").filter((h) => h.trim()) : undefined,
provider: provider || undefined,
+ target_node: targetNode || undefined,
}),
)
.then(() => {
@@ -94,7 +98,7 @@ const CreateRemoteDownload = () => {
.finally(() => {
setLoading(false);
});
- }, [target, url, path, fileName, username, password, headers, provider]);
+ }, [target, url, path, fileName, username, password, headers, provider, targetNode]);
return (
{
fullWidth
/>
+
{providers && providers.length > 1 && (
diff --git a/frontend/src/component/FileManager/Dialogs/ExtractArchive.tsx b/frontend/src/component/FileManager/Dialogs/ExtractArchive.tsx
index 3ddc5891..c1b106fb 100644
--- a/frontend/src/component/FileManager/Dialogs/ExtractArchive.tsx
+++ b/frontend/src/component/FileManager/Dialogs/ExtractArchive.tsx
@@ -9,6 +9,7 @@ import { fileExtension, getFileLinkedUri } from "../../../util";
import EncodingSelector, { defaultEncodingValue } from "../../Common/Form/EncodingSelector.tsx";
import { FileDisplayForm } from "../../Common/Form/FileDisplayForm.tsx";
import { PathSelectorForm } from "../../Common/Form/PathSelectorForm.tsx";
+import TargetNodeSelect, { useShowTargetNodeSelect } from "../../Common/Form/TargetNodeSelect.tsx";
import { ViewTaskAction } from "../../Common/Snackbar/snackbar.tsx";
import DraggableDialog from "../../Dialogs/DraggableDialog.tsx";
import Password from "../../Icons/Password.tsx";
@@ -26,6 +27,7 @@ const ExtractArchive = () => {
const [encoding, setEncoding] = useState(defaultEncodingValue);
const [password, setPassword] = useState("");
const [showPassword, setShowPassword] = useState(false);
+ const [targetNode, setTargetNode] = useState("");
const open = useAppSelector((state) => state.globalState.extractArchiveDialogOpen);
const target = useAppSelector((state) => state.globalState.extractArchiveDialogFile);
@@ -33,6 +35,7 @@ const ExtractArchive = () => {
const current = useAppSelector((state) => state.fileManager[FileManagerIndex.main].pure_path);
const mask = useAppSelector((state) => state.globalState.extractArchiveDialogMask);
const predefinedEncoding = useAppSelector((state) => state.globalState.extractArchiveDialogEncoding);
+ const showNodeSelect = useShowTargetNodeSelect();
useEffect(() => {
setEncoding(predefinedEncoding ?? defaultEncodingValue);
@@ -51,6 +54,7 @@ const ExtractArchive = () => {
useEffect(() => {
if (open) {
setPath(current ?? "");
+ setTargetNode("");
}
}, [open]);
@@ -71,6 +75,7 @@ const ExtractArchive = () => {
encoding: showEncodingOption && encoding != defaultEncodingValue ? encoding : undefined,
password: showPasswordOption && password ? password : undefined,
file_mask: mask ?? undefined,
+ target_node: targetNode || undefined,
}),
)
.then(() => {
@@ -84,7 +89,7 @@ const ExtractArchive = () => {
.finally(() => {
setLoading(false);
});
- }, [target, targets, encoding, path, showPasswordOption, showEncodingOption, password, mask]);
+ }, [target, targets, encoding, path, showPasswordOption, showEncodingOption, password, mask, targetNode]);
return (
{
>
+ {showNodeSelect && (
+
+
+
+ )}
{showPasswordOption && (
0 {
+ allowedSet := lo.SliceToMap(allowed, func(id int) (int, struct{}) { return id, struct{}{} })
+ filtered := lo.Filter(nodes, func(item *nodeItem, _ int) bool {
+ _, ok := allowedSet[item.node.ID()]
+ return ok
+ })
+ if len(filtered) == 0 {
+ return nil, fmt.Errorf("no allowed node found with capability %d: %w", capability, ErrNoAvailableNode)
+ }
+ nodes = filtered
+ }
+
var selected *nodeItem
if preferred > 0 {
@@ -198,6 +210,6 @@ func NewSlaveDummyNodePool(ctx context.Context, config conf.ConfigProvider, sett
func (s *slaveDummyNodePool) Upsert(ctx context.Context, node *ent.Node) {
}
-func (s *slaveDummyNodePool) Get(ctx context.Context, capability types.NodeCapability, preferred int) (Node, error) {
+func (s *slaveDummyNodePool) Get(ctx context.Context, capability types.NodeCapability, preferred int, allowed []int) (Node, error) {
return s.masterNode, nil
}
diff --git a/pkg/cluster/pool_test.go b/pkg/cluster/pool_test.go
index 5fc2cc89..f5e871b7 100644
--- a/pkg/cluster/pool_test.go
+++ b/pkg/cluster/pool_test.go
@@ -65,24 +65,60 @@ func TestWeightedNodePoolGet(t *testing.T) {
require.NoError(t, err)
// Capability with no matching nodes errors.
- _, err = pool.Get(ctx, types.NodeCapabilityCreateArchive, 0)
+ _, err = pool.Get(ctx, types.NodeCapabilityCreateArchive, 0, nil)
a.True(errors.Is(err, ErrNoAvailableNode))
// Preferred node wins regardless of weight.
- selected, err := pool.Get(ctx, types.NodeCapabilityRemoteDownload, n1.ID)
+ selected, err := pool.Get(ctx, types.NodeCapabilityRemoteDownload, n1.ID, nil)
require.NoError(t, err)
a.Equal(n1.ID, selected.ID())
// Without preference, the heavier node wins the first pick.
- selected, err = pool.Get(ctx, types.NodeCapabilityRemoteDownload, 0)
+ selected, err = pool.Get(ctx, types.NodeCapabilityRemoteDownload, 0, nil)
require.NoError(t, err)
a.Equal(n1.ID, selected.ID())
}
+func TestWeightedNodePoolGetAllowed(t *testing.T) {
+ a := assert.New(t)
+ client, pool, ctx := newPoolFixture(t)
+
+ n1 := client.Node.Create().
+ SetName("n1").SetType(node.TypeMaster).SetStatus(node.StatusActive).
+ SetServer("http://n1").SetCapabilities(capsOf(types.NodeCapabilityRemoteDownload)).
+ SetWeight(10).SaveX(ctx)
+ n2 := client.Node.Create().
+ SetName("n2").SetType(node.TypeMaster).SetStatus(node.StatusActive).
+ SetServer("http://n2").SetCapabilities(capsOf(types.NodeCapabilityRemoteDownload)).
+ SetWeight(1).SaveX(ctx)
+
+ pool, err := NewNodePool(ctx, logging.NewConsoleLogger(logging.LevelError), pool.(*weightedNodePool).conf, &stubSettings{}, inventory.NewNodeClient(client))
+ require.NoError(t, err)
+
+ // Allowed list narrows dispatch: heavier n1 excluded, n2 wins.
+ selected, err := pool.Get(ctx, types.NodeCapabilityRemoteDownload, 0, []int{n2.ID})
+ require.NoError(t, err)
+ a.Equal(n2.ID, selected.ID())
+
+ // Preferred node outside the allowed set falls back within the set.
+ selected, err = pool.Get(ctx, types.NodeCapabilityRemoteDownload, n1.ID, []int{n2.ID})
+ require.NoError(t, err)
+ a.Equal(n2.ID, selected.ID())
+
+ // Preferred node inside the allowed set wins.
+ selected, err = pool.Get(ctx, types.NodeCapabilityRemoteDownload, n1.ID, []int{n1.ID, n2.ID})
+ require.NoError(t, err)
+ a.Equal(n1.ID, selected.ID())
+
+ // Empty allowed intersection errors.
+ _, err = pool.Get(ctx, types.NodeCapabilityRemoteDownload, 0, []int{9999})
+ a.True(errors.Is(err, ErrNoAvailableNode))
+}
+
func TestWeightedNodePoolGetUnknownCapability(t *testing.T) {
a := assert.New(t)
_, pool, ctx := newPoolFixture(t)
- _, err := pool.Get(ctx, types.NodeCapabilityRemoteDownload, 0)
+ _, err := pool.Get(ctx, types.NodeCapabilityRemoteDownload, 0, nil)
a.True(errors.Is(err, ErrNoAvailableNode))
}
diff --git a/pkg/filemanager/workflows/archive.go b/pkg/filemanager/workflows/archive.go
index 2bef70b5..f98dcebb 100644
--- a/pkg/filemanager/workflows/archive.go
+++ b/pkg/filemanager/workflows/archive.go
@@ -75,11 +75,11 @@ func init() {
}
// NewCreateArchiveTask creates a new CreateArchiveTask
-func NewCreateArchiveTask(ctx context.Context, src []string, dst string) (queue.Task, error) {
+func NewCreateArchiveTask(ctx context.Context, src []string, dst string, sel NodeSelection) (queue.Task, error) {
state := &CreateArchiveTaskState{
Uris: src,
Dst: dst,
- NodeState: NodeState{},
+ NodeState: sel.state(),
}
stateBytes, err := json.Marshal(state)
if err != nil {
diff --git a/pkg/filemanager/workflows/extract.go b/pkg/filemanager/workflows/extract.go
index 8a72ddb2..840bd4a8 100644
--- a/pkg/filemanager/workflows/extract.go
+++ b/pkg/filemanager/workflows/extract.go
@@ -76,12 +76,12 @@ func init() {
// NewExtractArchiveTask creates a new ExtractArchiveTask. volumes optionally
// lists the URIs of all volumes of a multi-volume archive selected together.
-func NewExtractArchiveTask(ctx context.Context, src, dst, encoding, password string, mask []string, volumes []string) (queue.Task, error) {
+func NewExtractArchiveTask(ctx context.Context, src, dst, encoding, password string, mask []string, volumes []string, sel NodeSelection) (queue.Task, error) {
state := &ExtractArchiveTaskState{
Uri: src,
Dst: dst,
Encoding: encoding,
- NodeState: NodeState{},
+ NodeState: sel.state(),
Password: password,
FileMask: mask,
VolumeUris: volumes,
@@ -654,17 +654,17 @@ type (
}
SlaveExtractArchiveTaskState struct {
- FileName string `json:"file_name"`
- Entity *ent.Entity `json:"entity"`
- Policy *ent.StoragePolicy `json:"policy"`
- Encoding string `json:"encoding,omitempty"`
- Dst string `json:"dst,omitempty"`
- UserID int `json:"user_id"`
- TempPath string `json:"temp_path,omitempty"`
- TempZipFilePath string `json:"temp_zip_file_path,omitempty"`
- ProcessedCursor string `json:"processed_cursor,omitempty"`
- Password string `json:"password,omitempty"`
- FileMask []string `json:"file_mask,omitempty"`
+ FileName string `json:"file_name"`
+ Entity *ent.Entity `json:"entity"`
+ Policy *ent.StoragePolicy `json:"policy"`
+ Encoding string `json:"encoding,omitempty"`
+ Dst string `json:"dst,omitempty"`
+ UserID int `json:"user_id"`
+ TempPath string `json:"temp_path,omitempty"`
+ TempZipFilePath string `json:"temp_zip_file_path,omitempty"`
+ ProcessedCursor string `json:"processed_cursor,omitempty"`
+ Password string `json:"password,omitempty"`
+ FileMask []string `json:"file_mask,omitempty"`
Volumes map[string]*ent.Entity `json:"volumes,omitempty"`
}
)
@@ -698,7 +698,7 @@ func (m *SlaveExtractArchiveTask) Do(ctx context.Context) (task.Status, error) {
return task.StatusError, fmt.Errorf("failed to get node pool: %w", err)
}
- m.node, err = np.Get(ctx, types.NodeCapabilityNone, 0)
+ m.node, err = np.Get(ctx, types.NodeCapabilityNone, 0, nil)
if err != nil || !m.node.IsMaster() {
return task.StatusError, fmt.Errorf("failed to get master node: %w", err)
}
diff --git a/pkg/filemanager/workflows/remote_download.go b/pkg/filemanager/workflows/remote_download.go
index edcf2490..9d089af0 100644
--- a/pkg/filemanager/workflows/remote_download.go
+++ b/pkg/filemanager/workflows/remote_download.go
@@ -103,6 +103,8 @@ type RemoteDownloadTaskOption struct {
HTTPUsername string
HTTPPassword string
HTTPHeaders []string
+ // NodeSel carries the user/group node constraints for dispatch.
+ NodeSel NodeSelection
}
// NewRemoteDownloadTask creates a new RemoteDownloadTask
@@ -119,6 +121,7 @@ func NewRemoteDownloadTask(ctx context.Context, src string, srcFile, dst string,
state.HTTPUsername = opts.HTTPUsername
state.HTTPPassword = opts.HTTPPassword
state.HTTPHeaders = opts.HTTPHeaders
+ state.NodeState = opts.NodeSel.state()
}
stateBytes, err := json.Marshal(state)
if err != nil {
diff --git a/pkg/filemanager/workflows/upload.go b/pkg/filemanager/workflows/upload.go
index b50a0358..106dbe8b 100644
--- a/pkg/filemanager/workflows/upload.go
+++ b/pkg/filemanager/workflows/upload.go
@@ -76,7 +76,7 @@ func (t *SlaveUploadTask) Do(ctx context.Context) (task.Status, error) {
return task.StatusError, fmt.Errorf("failed to get node pool: %w", err)
}
- t.node, err = np.Get(ctx, types.NodeCapabilityNone, 0)
+ t.node, err = np.Get(ctx, types.NodeCapabilityNone, 0, nil)
if err != nil || !t.node.IsMaster() {
return task.StatusError, fmt.Errorf("failed to get master node: %w", err)
}
diff --git a/pkg/filemanager/workflows/workflows.go b/pkg/filemanager/workflows/workflows.go
index ddd35fbc..95c535b5 100644
--- a/pkg/filemanager/workflows/workflows.go
+++ b/pkg/filemanager/workflows/workflows.go
@@ -25,10 +25,25 @@ const (
type NodeState struct {
NodeID int `json:"node_id"`
+ // AllowedNodes restricts auto-dispatch to the group's node pool, captured
+ // at task creation so later group edits don't reroute queued tasks.
+ AllowedNodes []int `json:"allowed_nodes,omitempty"`
progress queue.Progresses
}
+// NodeSelection carries the caller-side node constraints for a new task:
+// TargetNodeID is an explicit user pick (0 = auto), AllowedNodes is the
+// group's eligible pool (empty = all).
+type NodeSelection struct {
+ TargetNodeID int
+ AllowedNodes []int
+}
+
+func (s NodeSelection) state() NodeState {
+ return NodeState{NodeID: s.TargetNodeID, AllowedNodes: s.AllowedNodes}
+}
+
// allocateNode allocates a node for the task.
func allocateNode(ctx context.Context, dep dependency.Dep, state *NodeState, capability types.NodeCapability) (cluster.Node, error) {
np, err := dep.NodePool(ctx)
@@ -36,7 +51,7 @@ func allocateNode(ctx context.Context, dep dependency.Dep, state *NodeState, cap
return nil, fmt.Errorf("failed to get node pool: %w", err)
}
- node, err := np.Get(ctx, capability, state.NodeID)
+ node, err := np.Get(ctx, capability, state.NodeID, state.AllowedNodes)
if err != nil {
return nil, fmt.Errorf("failed to get node: %w", err)
}
diff --git a/service/basic/site.go b/service/basic/site.go
index e444a3e2..2d58e4f5 100644
--- a/service/basic/site.go
+++ b/service/basic/site.go
@@ -1,12 +1,15 @@
package basic
import (
+ "slices"
"sort"
"strings"
"github.com/cloudreve/Cloudreve/v4/application/dependency"
+ "github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/inventory"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
+ "github.com/cloudreve/Cloudreve/v4/pkg/hashid"
"github.com/cloudreve/Cloudreve/v4/pkg/setting"
"github.com/cloudreve/Cloudreve/v4/pkg/thumb"
"github.com/cloudreve/Cloudreve/v4/service/user"
@@ -57,6 +60,12 @@ type SiteConfig struct {
// AbuseCaptcha controls whether the report-abuse dialog shows captcha.
AbuseCaptcha bool `json:"abuse_captcha,omitempty"`
+ // TaskNodes lists nodes the current user may target when creating tasks
+ // (remote download, archive ops); populated when the group allows node
+ // selection, filtered to the group's allowed pool.
+ TaskNodes []TaskNode `json:"task_nodes,omitempty"`
+ AllowSelectNode bool `json:"allow_select_node,omitempty"`
+
// Explorer section
Icons string `json:"icons,omitempty"`
EmojiPreset string `json:"emoji_preset,omitempty"`
@@ -94,6 +103,12 @@ type SiteConfig struct {
//AppForumLink string `json:"app_forum"`
}
+// TaskNode is the minimal public node descriptor for task targeting.
+type TaskNode struct {
+ ID string `json:"id"`
+ Name string `json:"name"`
+}
+
type (
GetSettingService struct {
Section string `uri:"section" binding:"required"`
@@ -217,6 +232,7 @@ func (s *GetSettingService) GetSiteConfig(c *gin.Context) (*SiteConfig, error) {
customNavItems := settings.CustomNavItems(c)
customHTML := settings.CustomHTML(c)
shareDefaults := settings.ShareDefaults(c)
+ taskNodes, allowSelect := taskNodesForUser(c, dep, u)
return &SiteConfig{
InstanceID: siteBasic.ID,
SiteName: siteBasic.Name,
@@ -238,9 +254,49 @@ func (s *GetSettingService) GetSiteConfig(c *gin.Context) (*SiteConfig, error) {
DefaultShareLinksInProfile: string(shareDefaults.LinksInProfile),
DownloadCDNRoutes: settings.DownloadCDNRoutes(c),
AbuseCaptcha: settings.AbuseCaptchaEnabled(c),
+ TaskNodes: taskNodes,
+ AllowSelectNode: allowSelect,
}, nil
}
+// taskNodesForUser returns the nodes a user may target for tasks: active
+// nodes with any task capability, intersected with the group's allowed pool.
+// Returns nil when the group disallows selection.
+func taskNodesForUser(c *gin.Context, dep dependency.Dep, u *ent.User) ([]TaskNode, bool) {
+ if u == nil || u.Edges.Group == nil || !u.Edges.Group.Settings.AllowSelectNode {
+ return nil, false
+ }
+
+ nodes, err := dep.NodeClient().ListActiveNodes(c, nil)
+ if err != nil {
+ return nil, true
+ }
+
+ allowed := u.Edges.Group.Settings.AllowedNodes
+ taskCaps := []types.NodeCapability{
+ types.NodeCapabilityCreateArchive,
+ types.NodeCapabilityExtractArchive,
+ types.NodeCapabilityRemoteDownload,
+ }
+ res := make([]TaskNode, 0, len(nodes))
+ for _, n := range nodes {
+ if len(allowed) > 0 && !slices.Contains(allowed, n.ID) {
+ continue
+ }
+ capable := false
+ for _, cap := range taskCaps {
+ if n.Capabilities != nil && n.Capabilities.Enabled(int(cap)) {
+ capable = true
+ break
+ }
+ }
+ if capable {
+ res = append(res, TaskNode{ID: hashid.EncodeNodeID(dep.HashIDEncoder(), n.ID), Name: n.Name})
+ }
+ }
+ return res, true
+}
+
const (
CaptchaSessionPrefix = "captcha_session_"
CaptchaTTL = 1800 // 30 minutes
diff --git a/service/explorer/workflows.go b/service/explorer/workflows.go
index 67864e68..93f50ef7 100644
--- a/service/explorer/workflows.go
+++ b/service/explorer/workflows.go
@@ -10,6 +10,7 @@ import (
"github.com/cloudreve/Cloudreve/v4/application/dependency"
"github.com/cloudreve/Cloudreve/v4/ent"
+ "github.com/cloudreve/Cloudreve/v4/ent/node"
"github.com/cloudreve/Cloudreve/v4/ent/task"
"github.com/cloudreve/Cloudreve/v4/inventory"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
@@ -25,6 +26,7 @@ import (
"github.com/gofrs/uuid"
"github.com/samber/lo"
"math"
+ "slices"
)
// ItemMoveService 处理多文件/目录移动
@@ -75,18 +77,53 @@ func init() {
type (
DownloadWorkflowService struct {
- Src []string `json:"src"`
- SrcFile string `json:"src_file"`
- Dst string `json:"dst" binding:"required"`
- FileName string `json:"file_name" binding:"omitempty,max=255"`
- Username string `json:"username" binding:"omitempty,max=255"`
- Password string `json:"password" binding:"omitempty,max=255"`
- Headers []string `json:"headers" binding:"omitempty,max=32,dive,max=2048"`
- Provider string `json:"provider" binding:"omitempty,max=64"`
+ Src []string `json:"src"`
+ SrcFile string `json:"src_file"`
+ Dst string `json:"dst" binding:"required"`
+ FileName string `json:"file_name" binding:"omitempty,max=255"`
+ Username string `json:"username" binding:"omitempty,max=255"`
+ Password string `json:"password" binding:"omitempty,max=255"`
+ Headers []string `json:"headers" binding:"omitempty,max=32,dive,max=2048"`
+ Provider string `json:"provider" binding:"omitempty,max=64"`
+ TargetNode string `json:"target_node" binding:"omitempty,max=64"`
}
CreateDownloadParamCtx struct{}
)
+// resolveNodeSelection validates the caller's target_node pick against group
+// constraints and returns the NodeSelection for task creation. An empty
+// target means auto dispatch; the group's allowed pool always applies.
+func resolveNodeSelection(c *gin.Context, dep dependency.Dep, target string, capability types.NodeCapability) (*workflows.NodeSelection, error) {
+ user := inventory.UserFromContext(c)
+ sel := &workflows.NodeSelection{AllowedNodes: user.Edges.Group.Settings.AllowedNodes}
+
+ if target == "" {
+ return sel, nil
+ }
+
+ if !user.Edges.Group.Settings.AllowSelectNode {
+ return nil, serializer.NewError(serializer.CodeGroupNotAllowed, "Group not allowed to select node", nil)
+ }
+
+ nodeID, err := dep.HashIDEncoder().Decode(target, hashid.NodeID)
+ if err != nil {
+ return nil, serializer.NewError(serializer.CodeParamErr, "Invalid target node", err)
+ }
+
+ if len(sel.AllowedNodes) > 0 && !slices.Contains(sel.AllowedNodes, nodeID) {
+ return nil, serializer.NewError(serializer.CodeParamErr, "Target node not allowed for this group", nil)
+ }
+
+ n, err := dep.NodeClient().GetNodeById(c, nodeID)
+ if err != nil || n.Status != node.StatusActive ||
+ n.Capabilities == nil || !n.Capabilities.Enabled(int(capability)) {
+ return nil, serializer.NewError(serializer.CodeParamErr, "Target node unavailable", err)
+ }
+
+ sel.TargetNodeID = nodeID
+ return sel, nil
+}
+
func (service *DownloadWorkflowService) CreateDownloadTask(c *gin.Context) ([]*TaskResponse, error) {
dep := dependency.FromContext(c)
user := inventory.UserFromContext(c)
@@ -169,6 +206,11 @@ func (service *DownloadWorkflowService) CreateDownloadTask(c *gin.Context) ([]*T
return nil, serializer.NewError(serializer.CodeParamErr, "Invalid downloader provider", nil)
}
+ nodeSel, err := resolveNodeSelection(c, dep, service.TargetNode, types.NodeCapabilityRemoteDownload)
+ if err != nil {
+ return nil, err
+ }
+
// Custom file name only applies to single-source tasks; HTTP credentials
// and headers only apply to plain HTTP(S) source URLs.
taskOpts := &workflows.RemoteDownloadTaskOption{
@@ -176,6 +218,7 @@ func (service *DownloadWorkflowService) CreateDownloadTask(c *gin.Context) ([]*T
HTTPUsername: service.Username,
HTTPPassword: service.Password,
HTTPHeaders: service.Headers,
+ NodeSel: *nodeSel,
}
if len(service.Src) <= 1 {
taskOpts.FileName = service.FileName
@@ -227,11 +270,12 @@ func (service *DownloadWorkflowService) CreateDownloadTask(c *gin.Context) ([]*T
type (
ArchiveWorkflowService struct {
- Src []string `json:"src" binding:"required"`
- Dst string `json:"dst" binding:"required"`
- Encoding string `json:"encoding"`
- Password string `json:"password"`
- FileMask []string `json:"file_mask"`
+ Src []string `json:"src" binding:"required"`
+ Dst string `json:"dst" binding:"required"`
+ Encoding string `json:"encoding"`
+ Password string `json:"password"`
+ FileMask []string `json:"file_mask"`
+ TargetNode string `json:"target_node" binding:"omitempty,max=64"`
}
CreateArchiveParamCtx struct{}
)
@@ -270,8 +314,13 @@ func (service *ArchiveWorkflowService) CreateExtractTask(c *gin.Context) (*TaskR
volumes = service.Src
}
+ nodeSel, err := resolveNodeSelection(c, dep, service.TargetNode, types.NodeCapabilityExtractArchive)
+ if err != nil {
+ return nil, err
+ }
+
// Create task
- t, err := workflows.NewExtractArchiveTask(c, src, service.Dst, service.Encoding, service.Password, service.FileMask, volumes)
+ t, err := workflows.NewExtractArchiveTask(c, src, service.Dst, service.Encoding, service.Password, service.FileMask, volumes, *nodeSel)
if err != nil {
return nil, serializer.NewError(serializer.CodeCreateTaskError, "Failed to create task", err)
}
@@ -316,8 +365,13 @@ func (service *ArchiveWorkflowService) CreateCompressTask(c *gin.Context) (*Task
}
m.OnUploadFailed(c, session)
+ nodeSel, err := resolveNodeSelection(c, dep, service.TargetNode, types.NodeCapabilityCreateArchive)
+ if err != nil {
+ return nil, err
+ }
+
// Create task
- t, err := workflows.NewCreateArchiveTask(c, service.Src, service.Dst)
+ t, err := workflows.NewCreateArchiveTask(c, service.Src, service.Dst, *nodeSel)
if err != nil {
return nil, serializer.NewError(serializer.CodeCreateTaskError, "Failed to create task", err)
}
diff --git a/service/explorer/workflows_nodeselect_test.go b/service/explorer/workflows_nodeselect_test.go
new file mode 100644
index 00000000..633f1477
--- /dev/null
+++ b/service/explorer/workflows_nodeselect_test.go
@@ -0,0 +1,127 @@
+package explorer
+
+import (
+ "context"
+ "fmt"
+ "net/http/httptest"
+ "testing"
+
+ "github.com/cloudreve/Cloudreve/v4/application/dependency"
+ "github.com/cloudreve/Cloudreve/v4/ent"
+ "github.com/cloudreve/Cloudreve/v4/ent/enttest"
+ entnode "github.com/cloudreve/Cloudreve/v4/ent/node"
+ entuser "github.com/cloudreve/Cloudreve/v4/ent/user"
+ "github.com/cloudreve/Cloudreve/v4/inventory"
+ "github.com/cloudreve/Cloudreve/v4/inventory/types"
+ "github.com/cloudreve/Cloudreve/v4/pkg/boolset"
+ "github.com/cloudreve/Cloudreve/v4/pkg/hashid"
+ "github.com/cloudreve/Cloudreve/v4/pkg/util"
+ "github.com/gin-gonic/gin"
+ "github.com/stretchr/testify/require"
+)
+
+// nodeSelDepStub exposes only the dependencies resolveNodeSelection touches.
+type nodeSelDepStub struct {
+ dependency.Dep
+ nodeClient inventory.NodeClient
+ hasher hashid.Encoder
+}
+
+func (d *nodeSelDepStub) NodeClient() inventory.NodeClient { return d.nodeClient }
+func (d *nodeSelDepStub) HashIDEncoder() hashid.Encoder { return d.hasher }
+
+func enableCaps(flags ...int) *boolset.BooleanSet {
+ b := boolset.BooleanSet(make([]byte, 4))
+ for _, f := range flags {
+ b[f/8] |= 1 << uint(f%8)
+ }
+ return &b
+}
+
+func TestResolveNodeSelection(t *testing.T) {
+ gin.SetMode(gin.TestMode)
+ client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared")
+ t.Cleanup(func() { require.NoError(t, client.Close()) })
+ ctx := context.Background()
+
+ mk := func(name string, status entnode.Status, caps *boolset.BooleanSet) *ent.Node {
+ return client.Node.Create().SetName(name).SetStatus(status).SetType(entnode.TypeSlave).
+ SetCapabilities(caps).SaveX(ctx)
+ }
+ downloadCapable := mk("dl", entnode.StatusActive, enableCaps(int(types.NodeCapabilityRemoteDownload)))
+ archiveOnly := mk("zip", entnode.StatusActive, enableCaps(int(types.NodeCapabilityCreateArchive)))
+ suspended := mk("off", entnode.StatusSuspended, enableCaps(int(types.NodeCapabilityRemoteDownload)))
+
+ seq := 0
+ newUser := func(settings *types.GroupSetting) *ent.User {
+ seq++
+ g := client.Group.Create().SetName("g").SetPermissions(&boolset.BooleanSet{}).
+ SetSettings(settings).SaveX(ctx)
+ u := client.User.Create().SetEmail(fmt.Sprintf("u%d@example.com", seq)).SetNick("u").SetStatus("active").SetGroup(g).SaveX(ctx)
+ return client.User.Query().WithGroup().Where(entuser.ID(u.ID)).OnlyX(ctx)
+ }
+
+ hasher, err := hashid.New("node-sel-test-salt")
+ require.NoError(t, err)
+ dep := &nodeSelDepStub{nodeClient: inventory.NewNodeClient(client), hasher: hasher}
+
+ newCtx := func(u *ent.User) *gin.Context {
+ engine := gin.New()
+ engine.ContextWithFallback = true
+ c := gin.CreateTestContextOnly(httptest.NewRecorder(), engine)
+ c.Request = httptest.NewRequest("POST", "/", nil)
+ util.WithValue(c, dependency.DepCtx{}, dep)
+ util.WithValue(c, inventory.UserCtx{}, u)
+ return c
+ }
+
+ t.Run("empty target returns allowed pool for auto dispatch", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowedNodes: []int{downloadCapable.ID}})
+ sel, err := resolveNodeSelection(newCtx(u), dep, "", types.NodeCapabilityRemoteDownload)
+ require.NoError(t, err)
+ require.Equal(t, 0, sel.TargetNodeID)
+ require.Equal(t, []int{downloadCapable.ID}, sel.AllowedNodes)
+ })
+
+ t.Run("explicit pick within allowed pool wins", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowedNodes: []int{downloadCapable.ID}, AllowSelectNode: true})
+ sel, err := resolveNodeSelection(newCtx(u), dep, hashid.EncodeNodeID(hasher, downloadCapable.ID),
+ types.NodeCapabilityRemoteDownload)
+ require.NoError(t, err)
+ require.Equal(t, downloadCapable.ID, sel.TargetNodeID)
+ })
+
+ t.Run("pick rejected when group disallows selection", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowSelectNode: false})
+ _, err := resolveNodeSelection(newCtx(u), dep, hashid.EncodeNodeID(hasher, downloadCapable.ID),
+ types.NodeCapabilityRemoteDownload)
+ require.Error(t, err)
+ })
+
+ t.Run("pick outside allowed pool is rejected", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowedNodes: []int{archiveOnly.ID}, AllowSelectNode: true})
+ _, err := resolveNodeSelection(newCtx(u), dep, hashid.EncodeNodeID(hasher, downloadCapable.ID),
+ types.NodeCapabilityRemoteDownload)
+ require.Error(t, err)
+ })
+
+ t.Run("node without the task capability is rejected", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowSelectNode: true})
+ _, err := resolveNodeSelection(newCtx(u), dep, hashid.EncodeNodeID(hasher, archiveOnly.ID),
+ types.NodeCapabilityRemoteDownload)
+ require.Error(t, err)
+ })
+
+ t.Run("suspended node is rejected", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowSelectNode: true})
+ _, err := resolveNodeSelection(newCtx(u), dep, hashid.EncodeNodeID(hasher, suspended.ID),
+ types.NodeCapabilityRemoteDownload)
+ require.Error(t, err)
+ })
+
+ t.Run("invalid node hashid is rejected", func(t *testing.T) {
+ u := newUser(&types.GroupSetting{AllowSelectNode: true})
+ _, err := resolveNodeSelection(newCtx(u), dep, "!!!", types.NodeCapabilityRemoteDownload)
+ require.Error(t, err)
+ })
+}