From 4d3ffd10e56768cccf0e1beac880bce882c212bf Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Thu, 11 Jan 2024 17:07:09 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E5=A4=9A=E5=8D=8F=E8=AE=AE?= =?UTF-8?q?=E4=BB=A3=E7=90=86=E7=9B=B8=E5=85=B3=E6=8E=A5=E5=8F=A3=20?= =?UTF-8?q?=E6=9B=B4=E6=96=B0=E5=AD=97=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/internal/http/cacah.go | 66 ------------------- client/internal/services/cacah.go | 57 ---------------- common/pkgs/cmd/create_package.go | 5 +- common/pkgs/cmd/update_package.go | 3 +- common/pkgs/db/cache.go | 8 +-- common/pkgs/db/model/model.go | 12 ---- common/pkgs/db/node.go | 13 ++-- common/pkgs/ioswitch/ops/grpc.go | 6 +- common/pkgs/ioswitch/plans/agent_plan.go | 7 +- common/pkgs/ioswitch/plans/plan_builder.go | 9 ++- common/pkgs/ioswitch/plans/plans.go | 4 +- .../pkgs/iterator/download_object_iterator.go | 6 +- common/pkgs/mq/coordinator/node.go | 9 ++- coordinator/internal/services/node.go | 4 +- .../event/check_package_redundancy.go | 3 +- scanner/internal/event/event_test.go | 21 +++--- .../tickevent/batch_all_agent_check_cache.go | 3 +- 17 files changed, 46 insertions(+), 190 deletions(-) delete mode 100644 client/internal/http/cacah.go delete mode 100644 client/internal/services/cacah.go diff --git a/client/internal/http/cacah.go b/client/internal/http/cacah.go deleted file mode 100644 index e1f1210..0000000 --- a/client/internal/http/cacah.go +++ /dev/null @@ -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 - } - } -} diff --git a/client/internal/services/cacah.go b/client/internal/services/cacah.go deleted file mode 100644 index ac37be9..0000000 --- a/client/internal/services/cacah.go +++ /dev/null @@ -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 -} diff --git a/common/pkgs/cmd/create_package.go b/common/pkgs/cmd/create_package.go index f374342..7a5dd77 100644 --- a/common/pkgs/cmd/create_package.go +++ b/common/pkgs/cmd/create_package.go @@ -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, diff --git a/common/pkgs/cmd/update_package.go b/common/pkgs/cmd/update_package.go index 50753fb..015644f 100644 --- a/common/pkgs/cmd/update_package.go +++ b/common/pkgs/cmd/update_package.go @@ -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, diff --git a/common/pkgs/db/cache.go b/common/pkgs/db/cache.go index d543592..ea59e8a 100644 --- a/common/pkgs/db/cache.go +++ b/common/pkgs/db/cache.go @@ -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"+ diff --git a/common/pkgs/db/model/model.go b/common/pkgs/db/model/model.go index f7b4eec..531d0e0 100644 --- a/common/pkgs/db/model/model.go +++ b/common/pkgs/db/model/model.go @@ -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"` diff --git a/common/pkgs/db/node.go b/common/pkgs/db/node.go index c600348..9460300 100644 --- a/common/pkgs/db/node.go +++ b/common/pkgs/db/node.go @@ -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 } diff --git a/common/pkgs/ioswitch/ops/grpc.go b/common/pkgs/ioswitch/ops/grpc.go index 104b258..b5d5cfe 100644 --- a/common/pkgs/ioswitch/ops/grpc.go +++ b/common/pkgs/ioswitch/ops/grpc.go @@ -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 { diff --git a/common/pkgs/ioswitch/plans/agent_plan.go b/common/pkgs/ioswitch/plans/agent_plan.go index f45f2e3..1f78a37 100644 --- a/common/pkgs/ioswitch/plans/agent_plan.go +++ b/common/pkgs/ioswitch/plans/agent_plan.go @@ -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(), diff --git a/common/pkgs/ioswitch/plans/plan_builder.go b/common/pkgs/ioswitch/plans/plan_builder.go index 771eb4a..6c95c1e 100644 --- a/common/pkgs/ioswitch/plans/plan_builder.go +++ b/common/pkgs/ioswitch/plans/plan_builder.go @@ -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 { diff --git a/common/pkgs/ioswitch/plans/plans.go b/common/pkgs/ioswitch/plans/plans.go index 3ea9e8e..60bafa8 100644 --- a/common/pkgs/ioswitch/plans/plans.go +++ b/common/pkgs/ioswitch/plans/plans.go @@ -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 } diff --git a/common/pkgs/iterator/download_object_iterator.go b/common/pkgs/iterator/download_object_iterator.go index 5cb1ef4..d04e660 100644 --- a/common/pkgs/iterator/download_object_iterator.go +++ b/common/pkgs/iterator/download_object_iterator.go @@ -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, diff --git a/common/pkgs/mq/coordinator/node.go b/common/pkgs/mq/coordinator/node.go index 7285f83..4eedea8 100644 --- a/common/pkgs/mq/coordinator/node.go +++ b/common/pkgs/mq/coordinator/node.go @@ -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, } diff --git a/coordinator/internal/services/node.go b/coordinator/internal/services/node.go index 0614264..c4c5683 100644 --- a/coordinator/internal/services/node.go +++ b/coordinator/internal/services/node.go @@ -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 diff --git a/scanner/internal/event/check_package_redundancy.go b/scanner/internal/event/check_package_redundancy.go index 9fa883d..357218a 100644 --- a/scanner/internal/event/check_package_redundancy.go +++ b/scanner/internal/event/check_package_redundancy.go @@ -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 } diff --git a/scanner/internal/event/event_test.go b/scanner/internal/event/event_test.go index 5ff44e1..64d6775 100644 --- a/scanner/internal/event/event_test.go +++ b/scanner/internal/event/event_test.go @@ -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}, diff --git a/scanner/internal/tickevent/batch_all_agent_check_cache.go b/scanner/internal/tickevent/batch_all_agent_check_cache.go index 2e21e00..dd08a28 100644 --- a/scanner/internal/tickevent/batch_all_agent_check_cache.go +++ b/scanner/internal/tickevent/batch_all_agent_check_cache.go @@ -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") }