| @@ -1,66 +0,0 @@ | |||
| package http | |||
| import ( | |||
| "net/http" | |||
| "time" | |||
| "github.com/gin-gonic/gin" | |||
| "gitlink.org.cn/cloudream/common/consts/errorcode" | |||
| "gitlink.org.cn/cloudream/common/pkgs/logger" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| ) | |||
| type CacheService struct { | |||
| *Server | |||
| } | |||
| func (s *Server) CacheSvc() *CacheService { | |||
| return &CacheService{ | |||
| Server: s, | |||
| } | |||
| } | |||
| type CacheMovePackageReq struct { | |||
| UserID *cdssdk.UserID `json:"userID" binding:"required"` | |||
| PackageID *cdssdk.PackageID `json:"packageID" binding:"required"` | |||
| NodeID *cdssdk.NodeID `json:"nodeID" binding:"required"` | |||
| } | |||
| type CacheMovePackageResp = cdssdk.CacheMovePackageResp | |||
| func (s *CacheService) MovePackage(ctx *gin.Context) { | |||
| log := logger.WithField("HTTP", "Cache.LoadPackage") | |||
| var req CacheMovePackageReq | |||
| if err := ctx.ShouldBindJSON(&req); err != nil { | |||
| log.Warnf("binding body: %s", err.Error()) | |||
| ctx.JSON(http.StatusBadRequest, Failed(errorcode.BadArgument, "missing argument or invalid argument")) | |||
| return | |||
| } | |||
| taskID, err := s.svc.CacheSvc().StartCacheMovePackage(*req.UserID, *req.PackageID, *req.NodeID) | |||
| if err != nil { | |||
| log.Warnf("start cache move package: %s", err.Error()) | |||
| ctx.JSON(http.StatusOK, Failed(errorcode.OperationFailed, "cache move package failed")) | |||
| return | |||
| } | |||
| for { | |||
| complete, err := s.svc.CacheSvc().WaitCacheMovePackage(*req.NodeID, taskID, time.Second*10) | |||
| if complete { | |||
| if err != nil { | |||
| log.Warnf("moving complete with: %s", err.Error()) | |||
| ctx.JSON(http.StatusOK, Failed(errorcode.OperationFailed, "cache move package failed")) | |||
| return | |||
| } | |||
| ctx.JSON(http.StatusOK, OK(CacheMovePackageResp{})) | |||
| return | |||
| } | |||
| if err != nil { | |||
| log.Warnf("wait moving: %s", err.Error()) | |||
| ctx.JSON(http.StatusOK, Failed(errorcode.OperationFailed, "cache move package failed")) | |||
| return | |||
| } | |||
| } | |||
| } | |||
| @@ -1,57 +0,0 @@ | |||
| package services | |||
| import ( | |||
| "fmt" | |||
| "time" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| stgglb "gitlink.org.cn/cloudream/storage/common/globals" | |||
| agtmq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/agent" | |||
| ) | |||
| type CacheService struct { | |||
| *Service | |||
| } | |||
| func (svc *Service) CacheSvc() *CacheService { | |||
| return &CacheService{Service: svc} | |||
| } | |||
| func (svc *CacheService) StartCacheMovePackage(userID cdssdk.UserID, packageID cdssdk.PackageID, nodeID cdssdk.NodeID) (string, error) { | |||
| agentCli, err := stgglb.AgentMQPool.Acquire(nodeID) | |||
| if err != nil { | |||
| return "", fmt.Errorf("new agent client: %w", err) | |||
| } | |||
| defer stgglb.AgentMQPool.Release(agentCli) | |||
| startResp, err := agentCli.StartCacheMovePackage(agtmq.NewStartCacheMovePackage(userID, packageID)) | |||
| if err != nil { | |||
| return "", fmt.Errorf("start cache move package: %w", err) | |||
| } | |||
| return startResp.TaskID, nil | |||
| } | |||
| func (svc *CacheService) WaitCacheMovePackage(nodeID cdssdk.NodeID, taskID string, waitTimeout time.Duration) (bool, error) { | |||
| agentCli, err := stgglb.AgentMQPool.Acquire(nodeID) | |||
| if err != nil { | |||
| return true, fmt.Errorf("new agent client: %w", err) | |||
| } | |||
| defer stgglb.AgentMQPool.Release(agentCli) | |||
| waitResp, err := agentCli.WaitCacheMovePackage(agtmq.NewWaitCacheMovePackage(taskID, waitTimeout.Milliseconds())) | |||
| if err != nil { | |||
| return true, fmt.Errorf("wait cache move package: %w", err) | |||
| } | |||
| if !waitResp.IsComplete { | |||
| return false, nil | |||
| } | |||
| if waitResp.Error != "" { | |||
| return true, fmt.Errorf("%s", waitResp.Error) | |||
| } | |||
| return true, nil | |||
| } | |||
| @@ -12,7 +12,6 @@ import ( | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| stgglb "gitlink.org.cn/cloudream/storage/common/globals" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/distlock/reqbuilder" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/iterator" | |||
| agtmq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/agent" | |||
| @@ -40,7 +39,7 @@ type ObjectUploadResult struct { | |||
| } | |||
| type UploadNodeInfo struct { | |||
| Node model.Node | |||
| Node cdssdk.Node | |||
| IsSameLocation bool | |||
| } | |||
| @@ -76,7 +75,7 @@ func (t *CreatePackage) Execute(ctx *UpdatePackageContext) (*CreatePackageResult | |||
| return nil, fmt.Errorf("getting user nodes: %w", err) | |||
| } | |||
| userNodes := lo.Map(getUserNodesResp.Nodes, func(node model.Node, index int) UploadNodeInfo { | |||
| userNodes := lo.Map(getUserNodesResp.Nodes, func(node cdssdk.Node, index int) UploadNodeInfo { | |||
| return UploadNodeInfo{ | |||
| Node: node, | |||
| IsSameLocation: node.LocationID == stgglb.Local.LocationID, | |||
| @@ -8,7 +8,6 @@ import ( | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| stgglb "gitlink.org.cn/cloudream/storage/common/globals" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/distlock/reqbuilder" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/iterator" | |||
| coormq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/coordinator" | |||
| @@ -50,7 +49,7 @@ func (t *UpdatePackage) Execute(ctx *UpdatePackageContext) (*UpdatePackageResult | |||
| return nil, fmt.Errorf("getting user nodes: %w", err) | |||
| } | |||
| userNodes := lo.Map(getUserNodesResp.Nodes, func(node model.Node, index int) UploadNodeInfo { | |||
| userNodes := lo.Map(getUserNodesResp.Nodes, func(node cdssdk.Node, index int) UploadNodeInfo { | |||
| return UploadNodeInfo{ | |||
| Node: node, | |||
| IsSameLocation: node.LocationID == stgglb.Local.LocationID, | |||
| @@ -74,8 +74,8 @@ func (*CacheDB) NodeBatchDelete(ctx SQLContext, nodeID cdssdk.NodeID, fileHashes | |||
| } | |||
| // GetCachingFileNodes 查找缓存了指定文件的节点 | |||
| func (*CacheDB) GetCachingFileNodes(ctx SQLContext, fileHash string) ([]model.Node, error) { | |||
| var x []model.Node | |||
| func (*CacheDB) GetCachingFileNodes(ctx SQLContext, fileHash string) ([]cdssdk.Node, error) { | |||
| var x []cdssdk.Node | |||
| err := sqlx.Select(ctx, &x, | |||
| "select Node.* from Cache, Node where Cache.FileHash=? and Cache.NodeID = Node.NodeID", fileHash) | |||
| return x, err | |||
| @@ -88,8 +88,8 @@ func (*CacheDB) DeleteNodeAll(ctx SQLContext, nodeID cdssdk.NodeID) error { | |||
| } | |||
| // FindCachingFileUserNodes 在缓存表中查询指定数据所在的节点 | |||
| func (*CacheDB) FindCachingFileUserNodes(ctx SQLContext, userID cdssdk.NodeID, fileHash string) ([]model.Node, error) { | |||
| var x []model.Node | |||
| func (*CacheDB) FindCachingFileUserNodes(ctx SQLContext, userID cdssdk.NodeID, fileHash string) ([]cdssdk.Node, error) { | |||
| var x []cdssdk.Node | |||
| err := sqlx.Select(ctx, &x, | |||
| "select Node.* from Cache, UserNode, Node where"+ | |||
| " Cache.FileHash=? and Cache.NodeID = UserNode.NodeID and"+ | |||
| @@ -12,18 +12,6 @@ import ( | |||
| // TODO 可以考虑逐步迁移到cdssdk中。迁移思路:数据对象应该包含的字段都迁移到cdssdk中,内部使用的一些特殊字段则留在这里 | |||
| type Node struct { | |||
| NodeID cdssdk.NodeID `db:"NodeID" json:"nodeID"` | |||
| Name string `db:"Name" json:"name"` | |||
| LocalIP string `db:"LocalIP" json:"localIP"` | |||
| ExternalIP string `db:"ExternalIP" json:"externalIP"` | |||
| LocalGRPCPort int `db:"LocalGRPCPort" json:"localGRPCPort"` | |||
| ExternalGRPCPort int `db:"ExternalGRPCPort" json:"externalGRPCPort"` | |||
| LocationID cdssdk.LocationID `db:"LocationID" json:"locationID"` | |||
| State string `db:"State" json:"state"` | |||
| LastReportTime *time.Time `db:"LastReportTime" json:"lastReportTime"` | |||
| } | |||
| type Storage struct { | |||
| StorageID cdssdk.StorageID `db:"StorageID" json:"storageID"` | |||
| Name string `db:"Name" json:"name"` | |||
| @@ -5,7 +5,6 @@ import ( | |||
| "github.com/jmoiron/sqlx" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| ) | |||
| type NodeDB struct { | |||
| @@ -16,21 +15,21 @@ func (db *DB) Node() *NodeDB { | |||
| return &NodeDB{DB: db} | |||
| } | |||
| func (db *NodeDB) GetByID(ctx SQLContext, nodeID cdssdk.NodeID) (model.Node, error) { | |||
| var ret model.Node | |||
| func (db *NodeDB) GetByID(ctx SQLContext, nodeID cdssdk.NodeID) (cdssdk.Node, error) { | |||
| var ret cdssdk.Node | |||
| err := sqlx.Get(ctx, &ret, "select * from Node where NodeID = ?", nodeID) | |||
| return ret, err | |||
| } | |||
| func (db *NodeDB) GetAllNodes(ctx SQLContext) ([]model.Node, error) { | |||
| var ret []model.Node | |||
| func (db *NodeDB) GetAllNodes(ctx SQLContext) ([]cdssdk.Node, error) { | |||
| var ret []cdssdk.Node | |||
| err := sqlx.Select(ctx, &ret, "select * from Node") | |||
| return ret, err | |||
| } | |||
| // GetUserNodes 根据用户id查询可用node | |||
| func (db *NodeDB) GetUserNodes(ctx SQLContext, userID cdssdk.UserID) ([]model.Node, error) { | |||
| var nodes []model.Node | |||
| func (db *NodeDB) GetUserNodes(ctx SQLContext, userID cdssdk.UserID) ([]cdssdk.Node, error) { | |||
| var nodes []cdssdk.Node | |||
| err := sqlx.Select(ctx, &nodes, "select Node.* from UserNode, Node where UserNode.NodeID = Node.NodeID and UserNode.UserID=?", userID) | |||
| return nodes, err | |||
| } | |||
| @@ -7,16 +7,16 @@ import ( | |||
| "gitlink.org.cn/cloudream/common/pkgs/future" | |||
| "gitlink.org.cn/cloudream/common/pkgs/logger" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| myio "gitlink.org.cn/cloudream/common/utils/io" | |||
| stgglb "gitlink.org.cn/cloudream/storage/common/globals" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/ioswitch" | |||
| ) | |||
| type GRPCSend struct { | |||
| LocalID ioswitch.StreamID `json:"localID"` | |||
| RemoteID ioswitch.StreamID `json:"remoteID"` | |||
| Node model.Node `json:"node"` | |||
| Node cdssdk.Node `json:"node"` | |||
| } | |||
| func (o *GRPCSend) Execute(sw *ioswitch.Switch, planID ioswitch.PlanID) error { | |||
| @@ -49,7 +49,7 @@ func (o *GRPCSend) Execute(sw *ioswitch.Switch, planID ioswitch.PlanID) error { | |||
| type GRPCFetch struct { | |||
| RemoteID ioswitch.StreamID `json:"remoteID"` | |||
| LocalID ioswitch.StreamID `json:"localID"` | |||
| Node model.Node `json:"node"` | |||
| Node cdssdk.Node `json:"node"` | |||
| } | |||
| func (o *GRPCFetch) Execute(sw *ioswitch.Switch, planID ioswitch.PlanID) error { | |||
| @@ -2,14 +2,13 @@ package plans | |||
| import ( | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/ioswitch" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/ioswitch/ops" | |||
| ) | |||
| type AgentPlanBuilder struct { | |||
| owner *PlanBuilder | |||
| node model.Node | |||
| node cdssdk.Node | |||
| ops []ioswitch.Op | |||
| } | |||
| @@ -30,7 +29,7 @@ func (b *AgentPlanBuilder) Build(planID ioswitch.PlanID) (AgentPlan, error) { | |||
| }, nil | |||
| } | |||
| func (b *AgentPlanBuilder) GRCPFetch(node model.Node, str *AgentStream) *AgentStream { | |||
| func (b *AgentPlanBuilder) GRCPFetch(node cdssdk.Node, str *AgentStream) *AgentStream { | |||
| agtStr := &AgentStream{ | |||
| owner: b, | |||
| info: b.owner.newStream(), | |||
| @@ -45,7 +44,7 @@ func (b *AgentPlanBuilder) GRCPFetch(node model.Node, str *AgentStream) *AgentSt | |||
| return agtStr | |||
| } | |||
| func (s *AgentStream) GRPCSend(node model.Node) *AgentStream { | |||
| func (s *AgentStream) GRPCSend(node cdssdk.Node) *AgentStream { | |||
| agtStr := &AgentStream{ | |||
| owner: s.owner.owner.AtAgent(node), | |||
| info: s.owner.owner.newStream(), | |||
| @@ -5,7 +5,6 @@ import ( | |||
| "github.com/google/uuid" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/ioswitch" | |||
| ) | |||
| @@ -60,7 +59,7 @@ func (b *PlanBuilder) FromExecutor() *FromExecutorStream { | |||
| } | |||
| } | |||
| func (b *PlanBuilder) AtAgent(node model.Node) *AgentPlanBuilder { | |||
| func (b *PlanBuilder) AtAgent(node cdssdk.Node) *AgentPlanBuilder { | |||
| agtPlan, ok := b.agentPlans[node.NodeID] | |||
| if !ok { | |||
| agtPlan = &AgentPlanBuilder{ | |||
| @@ -76,10 +75,10 @@ func (b *PlanBuilder) AtAgent(node model.Node) *AgentPlanBuilder { | |||
| type FromExecutorStream struct { | |||
| owner *PlanBuilder | |||
| info *StreamInfo | |||
| toNode *model.Node | |||
| toNode *cdssdk.Node | |||
| } | |||
| func (s *FromExecutorStream) ToNode(node model.Node) *AgentStream { | |||
| func (s *FromExecutorStream) ToNode(node cdssdk.Node) *AgentStream { | |||
| s.toNode = &node | |||
| return &AgentStream{ | |||
| owner: s.owner.AtAgent(node), | |||
| @@ -89,7 +88,7 @@ func (s *FromExecutorStream) ToNode(node model.Node) *AgentStream { | |||
| type ToExecutorStream struct { | |||
| info *StreamInfo | |||
| fromNode *model.Node | |||
| fromNode *cdssdk.Node | |||
| } | |||
| type MultiStream struct { | |||
| @@ -1,12 +1,12 @@ | |||
| package plans | |||
| import ( | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/ioswitch" | |||
| ) | |||
| type AgentPlan struct { | |||
| Node model.Node | |||
| Node cdssdk.Node | |||
| Plan ioswitch.Plan | |||
| } | |||
| @@ -28,7 +28,7 @@ type IterDownloadingObject struct { | |||
| } | |||
| type DownloadNodeInfo struct { | |||
| Node model.Node | |||
| Node cdssdk.Node | |||
| IsSameLocation bool | |||
| } | |||
| @@ -141,7 +141,7 @@ func (iter *DownloadObjectIterator) downloadNoneOrRepObject(coorCli *coormq.Clie | |||
| continue | |||
| } | |||
| downloadNodes := lo.Map(getNodesResp.Nodes, func(node model.Node, index int) DownloadNodeInfo { | |||
| downloadNodes := lo.Map(getNodesResp.Nodes, func(node cdssdk.Node, index int) DownloadNodeInfo { | |||
| return DownloadNodeInfo{ | |||
| Node: node, | |||
| IsSameLocation: node.LocationID == stgglb.Local.LocationID, | |||
| @@ -186,7 +186,7 @@ func (iter *DownloadObjectIterator) downloadECObject(coorCli *coormq.Client, ctx | |||
| continue | |||
| } | |||
| downloadNodes := lo.Map(getNodesResp.Nodes, func(node model.Node, index int) DownloadNodeInfo { | |||
| downloadNodes := lo.Map(getNodesResp.Nodes, func(node cdssdk.Node, index int) DownloadNodeInfo { | |||
| return DownloadNodeInfo{ | |||
| Node: node, | |||
| IsSameLocation: node.LocationID == stgglb.Local.LocationID, | |||
| @@ -3,7 +3,6 @@ package coordinator | |||
| import ( | |||
| "gitlink.org.cn/cloudream/common/pkgs/mq" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| ) | |||
| type NodeService interface { | |||
| @@ -21,7 +20,7 @@ type GetUserNodes struct { | |||
| } | |||
| type GetUserNodesResp struct { | |||
| mq.MessageBodyBase | |||
| Nodes []model.Node `json:"nodes"` | |||
| Nodes []cdssdk.Node `json:"nodes"` | |||
| } | |||
| func NewGetUserNodes(userID cdssdk.UserID) *GetUserNodes { | |||
| @@ -29,7 +28,7 @@ func NewGetUserNodes(userID cdssdk.UserID) *GetUserNodes { | |||
| UserID: userID, | |||
| } | |||
| } | |||
| func NewGetUserNodesResp(nodes []model.Node) *GetUserNodesResp { | |||
| func NewGetUserNodesResp(nodes []cdssdk.Node) *GetUserNodesResp { | |||
| return &GetUserNodesResp{ | |||
| Nodes: nodes, | |||
| } | |||
| @@ -47,7 +46,7 @@ type GetNodes struct { | |||
| } | |||
| type GetNodesResp struct { | |||
| mq.MessageBodyBase | |||
| Nodes []model.Node `json:"nodes"` | |||
| Nodes []cdssdk.Node `json:"nodes"` | |||
| } | |||
| func NewGetNodes(nodeIDs []cdssdk.NodeID) *GetNodes { | |||
| @@ -55,7 +54,7 @@ func NewGetNodes(nodeIDs []cdssdk.NodeID) *GetNodes { | |||
| NodeIDs: nodeIDs, | |||
| } | |||
| } | |||
| func NewGetNodesResp(nodes []model.Node) *GetNodesResp { | |||
| func NewGetNodesResp(nodes []cdssdk.Node) *GetNodesResp { | |||
| return &GetNodesResp{ | |||
| Nodes: nodes, | |||
| } | |||
| @@ -4,7 +4,7 @@ import ( | |||
| "gitlink.org.cn/cloudream/common/consts/errorcode" | |||
| "gitlink.org.cn/cloudream/common/pkgs/logger" | |||
| "gitlink.org.cn/cloudream/common/pkgs/mq" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| coormq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/coordinator" | |||
| ) | |||
| @@ -20,7 +20,7 @@ func (svc *Service) GetUserNodes(msg *coormq.GetUserNodes) (*coormq.GetUserNodes | |||
| } | |||
| func (svc *Service) GetNodes(msg *coormq.GetNodes) (*coormq.GetNodesResp, *mq.CodeMessage) { | |||
| var nodes []model.Node | |||
| var nodes []cdssdk.Node | |||
| if msg.NodeIDs == nil { | |||
| var err error | |||
| @@ -11,7 +11,6 @@ import ( | |||
| "gitlink.org.cn/cloudream/common/utils/sort" | |||
| stgglb "gitlink.org.cn/cloudream/storage/common/globals" | |||
| stgmod "gitlink.org.cn/cloudream/storage/common/models" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/distlock/reqbuilder" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/ioswitch/plans" | |||
| agtmq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/agent" | |||
| @@ -36,7 +35,7 @@ func NewCheckPackageRedundancy(evt *scevt.CheckPackageRedundancy) *CheckPackageR | |||
| } | |||
| type NodeLoadInfo struct { | |||
| Node model.Node | |||
| Node cdssdk.Node | |||
| LoadsRecentMonth int | |||
| LoadsRecentYear int | |||
| } | |||
| @@ -6,7 +6,6 @@ import ( | |||
| "github.com/samber/lo" | |||
| . "github.com/smartystreets/goconvey/convey" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| ) | |||
| func Test_chooseSoManyNodes(t *testing.T) { | |||
| @@ -19,8 +18,8 @@ func Test_chooseSoManyNodes(t *testing.T) { | |||
| { | |||
| title: "节点数量充足", | |||
| allNodes: []*NodeLoadInfo{ | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(2)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(2)}}, | |||
| }, | |||
| count: 2, | |||
| expectedNodeIDs: []cdssdk.NodeID{1, 2}, | |||
| @@ -28,9 +27,9 @@ func Test_chooseSoManyNodes(t *testing.T) { | |||
| { | |||
| title: "节点数量超过", | |||
| allNodes: []*NodeLoadInfo{ | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(2)}}, | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(3)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(2)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(3)}}, | |||
| }, | |||
| count: 2, | |||
| expectedNodeIDs: []cdssdk.NodeID{1, 2}, | |||
| @@ -38,7 +37,7 @@ func Test_chooseSoManyNodes(t *testing.T) { | |||
| { | |||
| title: "只有一个节点,节点数量不够", | |||
| allNodes: []*NodeLoadInfo{ | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| }, | |||
| count: 3, | |||
| expectedNodeIDs: []cdssdk.NodeID{1, 1, 1}, | |||
| @@ -46,8 +45,8 @@ func Test_chooseSoManyNodes(t *testing.T) { | |||
| { | |||
| title: "多个同地区节点,节点数量不够", | |||
| allNodes: []*NodeLoadInfo{ | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(2)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(1)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(2)}}, | |||
| }, | |||
| count: 5, | |||
| expectedNodeIDs: []cdssdk.NodeID{1, 1, 1, 2, 2}, | |||
| @@ -55,8 +54,8 @@ func Test_chooseSoManyNodes(t *testing.T) { | |||
| { | |||
| title: "节点数量不够,且在不同地区", | |||
| allNodes: []*NodeLoadInfo{ | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(1), LocationID: cdssdk.LocationID(1)}}, | |||
| {Node: model.Node{NodeID: cdssdk.NodeID(2), LocationID: cdssdk.LocationID(2)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(1), LocationID: cdssdk.LocationID(1)}}, | |||
| {Node: cdssdk.Node{NodeID: cdssdk.NodeID(2), LocationID: cdssdk.LocationID(2)}}, | |||
| }, | |||
| count: 5, | |||
| expectedNodeIDs: []cdssdk.NodeID{1, 2, 1, 2, 1}, | |||
| @@ -4,7 +4,6 @@ import ( | |||
| "github.com/samber/lo" | |||
| "gitlink.org.cn/cloudream/common/pkgs/logger" | |||
| cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" | |||
| "gitlink.org.cn/cloudream/storage/common/pkgs/db/model" | |||
| scevt "gitlink.org.cn/cloudream/storage/common/pkgs/mq/scanner/event" | |||
| "gitlink.org.cn/cloudream/storage/scanner/internal/event" | |||
| ) | |||
| @@ -31,7 +30,7 @@ func (e *BatchAllAgentCheckCache) Execute(ctx ExecuteContext) { | |||
| return | |||
| } | |||
| e.nodeIDs = lo.Map(nodes, func(node model.Node, index int) cdssdk.NodeID { return node.NodeID }) | |||
| e.nodeIDs = lo.Map(nodes, func(node cdssdk.Node, index int) cdssdk.NodeID { return node.NodeID }) | |||
| log.Debugf("new check start, get all nodes") | |||
| } | |||