Merge branch 'split_coor'

This commit is contained in:
Sydonian 2023-08-14 11:31:39 +08:00
commit fa38c4fcdc
11 changed files with 877 additions and 0 deletions

33
go.mod Normal file
View File

@ -0,0 +1,33 @@
module gitlink.org.cn/cloudream/storage-coordinator
go 1.20
require (
github.com/jmoiron/sqlx v1.3.5
github.com/samber/lo v1.38.1
gitlink.org.cn/cloudream/common v0.0.0
gitlink.org.cn/cloudream/storage-common v0.0.0
)
require (
github.com/antonfisher/nested-logrus-formatter v1.3.1 // indirect
github.com/go-sql-driver/mysql v1.7.1 // indirect
github.com/google/uuid v1.3.0 // indirect
github.com/hashicorp/errwrap v1.1.0 // indirect
github.com/hashicorp/go-multierror v1.1.1 // indirect
github.com/imdario/mergo v0.3.15 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/mitchellh/mapstructure v1.5.0 // indirect
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421 // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/sirupsen/logrus v1.9.2 // indirect
github.com/streadway/amqp v1.1.0 // indirect
github.com/zyedidia/generic v1.2.1 // indirect
golang.org/x/exp v0.0.0-20230519143937-03e91628a987 // indirect
golang.org/x/sys v0.7.0 // indirect
)
// go mod tidy时需要将下面几行取消注释
replace gitlink.org.cn/cloudream/common => ../../common
replace gitlink.org.cn/cloudream/storage-common => ../storage-common

59
go.sum Normal file
View File

@ -0,0 +1,59 @@
github.com/antonfisher/nested-logrus-formatter v1.3.1 h1:NFJIr+pzwv5QLHTPyKz9UMEoHck02Q9L0FP13b/xSbQ=
github.com/antonfisher/nested-logrus-formatter v1.3.1/go.mod h1:6WTfyWFkBc9+zyBaKIqRrg/KwMqBbodBjgbHjDz7zjA=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/go-sql-driver/mysql v1.6.0/go.mod h1:DCzpHaOWr8IXmIStZouvnhqoel9Qv2LBy8hT2VhHyBg=
github.com/go-sql-driver/mysql v1.7.1 h1:lUIinVbN1DY0xBg0eMOzmmtGoHwWBbvnWubQUrtU8EI=
github.com/go-sql-driver/mysql v1.7.1/go.mod h1:OXbVy3sEdcQ2Doequ6Z5BW6fXNQTmx+9S1MCJN5yJMI=
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/uuid v1.3.0 h1:t6JiXgmwXMjEs8VusXIJk2BXHsn+wx8BZdTaoZ5fu7I=
github.com/google/uuid v1.3.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gopherjs/gopherjs v1.17.2 h1:fQnZVsXk8uxXIStYb0N4bGk7jeyTalG/wsZjQ25dO0g=
github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
github.com/hashicorp/errwrap v1.1.0 h1:OxrOeh75EUXMY8TBjag2fzXGZ40LB6IKw45YeGUDY2I=
github.com/hashicorp/errwrap v1.1.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+lD48awMYo=
github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM=
github.com/imdario/mergo v0.3.15 h1:M8XP7IuFNsqUx6VPK2P9OSmsYsI/YFaGil0uD21V3dM=
github.com/imdario/mergo v0.3.15/go.mod h1:WBLT9ZmE3lPoWsEzCh9LPo3TiwVN+ZKEjmz+hD27ysY=
github.com/jmoiron/sqlx v1.3.5 h1:vFFPA71p1o5gAeqtEAwLU4dnX2napprKtHr7PYIcN3g=
github.com/jmoiron/sqlx v1.3.5/go.mod h1:nRVWtLre0KfCLJvgxzCsLVMogSvQ1zNJtpYr2Ccp0mQ=
github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM=
github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo=
github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7C0MuV77Wo=
github.com/lib/pq v1.2.0 h1:LXpIM/LZ5xGFhOpXAQUIMM1HdyqzVYM13zNdjCEEcA0=
github.com/lib/pq v1.2.0/go.mod h1:5WUZQaWbwv1U+lTReE5YruASi9Al49XbQIvNi/34Woo=
github.com/mattn/go-sqlite3 v1.14.6 h1:dNPt6NO46WmLVt2DLNpwczCmdV5boIZ6g/tlDrlRUbg=
github.com/mattn/go-sqlite3 v1.14.6/go.mod h1:NyWgC/yNuGj7Q9rpYnZvas74GogHl5/Z4A/KQRfk6bU=
github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyuac5Z2hdY=
github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo=
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421 h1:ZqeYNhU3OHLH3mGKHDcjJRFFRrJa6eAM5H+CtDdOsPc=
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M=
github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/samber/lo v1.38.1 h1:j2XEAqXKb09Am4ebOg31SpvzUTTs6EN3VfgeLUhPdXM=
github.com/samber/lo v1.38.1/go.mod h1:+m/ZKRl6ClXCE2Lgf3MsQlWfh4bn1bz6CXEOxnEXnEA=
github.com/sirupsen/logrus v1.9.2 h1:oxx1eChJGI6Uks2ZC4W1zpLlVgqB8ner4EuQwV4Ik1Y=
github.com/sirupsen/logrus v1.9.2/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
github.com/smartystreets/assertions v1.13.1 h1:Ef7KhSmjZcK6AVf9YbJdvPYG9avaF0ZxudX+ThRdWfU=
github.com/smartystreets/goconvey v1.8.0 h1:Oi49ha/2MURE0WexF052Z0m+BNSGirfjg5RL+JXWq3w=
github.com/streadway/amqp v1.1.0 h1:py12iX8XSyI7aN/3dUT8DFIDJazNJsVJdxNVEpnQTZM=
github.com/streadway/amqp v1.1.0/go.mod h1:WYSrTEYHOXHd0nwFeUXAe2G2hRnQT+deZJJf88uS9Bg=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/zyedidia/generic v1.2.1 h1:Zv5KS/N2m0XZZiuLS82qheRG4X1o5gsWreGb0hR7XDc=
github.com/zyedidia/generic v1.2.1/go.mod h1:ly2RBz4mnz1yeuVbQA/VFwGjK3mnHGRj1JuoG336Bis=
golang.org/x/exp v0.0.0-20230519143937-03e91628a987 h1:3xJIFvzUFbu4ls0BTBYcgbCGhA63eAOEMxIHugyXJqA=
golang.org/x/exp v0.0.0-20230519143937-03e91628a987/go.mod h1:V1LtkGg67GoY2N1AnLN78QLrzxkLyJw7RJb1gzOOz9w=
golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.7.0 h1:3jlCCIQZPdOYu1h8BkNvLz8Kgwtae2cagcG/VamtZRU=
golang.org/x/sys v0.7.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=

24
internal/config/config.go Normal file
View File

@ -0,0 +1,24 @@
package config
import (
log "gitlink.org.cn/cloudream/common/pkg/logger"
c "gitlink.org.cn/cloudream/common/utils/config"
db "gitlink.org.cn/cloudream/storage-common/pkgs/db/config"
racfg "gitlink.org.cn/cloudream/storage-common/pkgs/mq/config"
)
type Config struct {
Logger log.Config `json:"logger"`
DB db.Config `json:"db"`
RabbitMQ racfg.Config `json:"rabbitMQ"`
}
var cfg Config
func Init() error {
return c.DefaultLoad("coordinator", &cfg)
}
func Cfg() *Config {
return &cfg
}

View File

@ -0,0 +1,24 @@
package services
import (
coormsg "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message/coordinator"
)
func (service *Service) TempCacheReport(msg *coormsg.TempCacheReport) {
service.db.BatchInsertOrUpdateCache(msg.Hashes, msg.NodeID)
}
func (service *Service) AgentStatusReport(msg *coormsg.AgentStatusReport) {
//jh根据command中的Ip插入节点延迟表和节点表的NodeStatus
//根据command中的Ip插入节点延迟表
// TODO
/*
ips := utils.GetAgentIps()
Insert_NodeDelay(msg.Body.IP, ips, msg.Body.AgentDelay)
//从配置表里读取节点地域NodeLocation
//插入节点表的NodeStatus
Insert_Node(msg.Body.IP, msg.Body.IP, msg.Body.IPFSStatus, msg.Body.LocalDirStatus)
*/
}

View File

@ -0,0 +1,74 @@
package services
import (
"database/sql"
"github.com/jmoiron/sqlx"
"gitlink.org.cn/cloudream/common/consts/errorcode"
log "gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkg/mq"
"gitlink.org.cn/cloudream/storage-common/pkgs/db/model"
coormsg "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message/coordinator"
)
func (svc *Service) GetBucket(userID int, bucketID int) (model.Bucket, error) {
// TODO
panic("not implement yet")
}
func (svc *Service) GetUserBuckets(msg *coormsg.GetUserBuckets) (*coormsg.GetUserBucketsResp, *mq.CodeMessage) {
buckets, err := svc.db.Bucket().GetUserBuckets(svc.db.SQLCtx(), msg.UserID)
if err != nil {
log.WithField("UserID", msg.UserID).
Warnf("get user buckets failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.GetUserBucketsResp](errorcode.OperationFailed, "get all buckets failed")
}
return mq.ReplyOK(coormsg.NewGetUserBucketsResp(buckets))
}
func (svc *Service) GetBucketObjects(msg *coormsg.GetBucketObjects) (*coormsg.GetBucketObjectsResp, *mq.CodeMessage) {
objects, err := svc.db.Object().GetBucketObjects(svc.db.SQLCtx(), msg.UserID, msg.BucketID)
if err != nil {
log.WithField("UserID", msg.UserID).
WithField("BucketID", msg.BucketID).
Warnf("get bucket objects failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.GetBucketObjectsResp](errorcode.OperationFailed, "get bucket objects failed")
}
return mq.ReplyOK(coormsg.NewGetBucketObjectsResp(objects))
}
func (svc *Service) CreateBucket(msg *coormsg.CreateBucket) (*coormsg.CreateBucketResp, *mq.CodeMessage) {
var bucketID int64
var err error
svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
// 这里用的是外部的err
bucketID, err = svc.db.Bucket().Create(tx, msg.UserID, msg.BucketName)
return err
})
if err != nil {
log.WithField("UserID", msg.UserID).
WithField("BucketName", msg.BucketName).
Warnf("create bucket failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.CreateBucketResp](errorcode.OperationFailed, "create bucket failed")
}
return mq.ReplyOK(coormsg.NewCreateBucketResp(bucketID))
}
func (svc *Service) DeleteBucket(msg *coormsg.DeleteBucket) (*coormsg.DeleteBucketResp, *mq.CodeMessage) {
err := svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
return svc.db.Bucket().Delete(tx, msg.BucketID)
})
if err != nil {
log.WithField("UserID", msg.UserID).
WithField("BucketID", msg.BucketID).
Warnf("delete bucket failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.DeleteBucketResp](errorcode.OperationFailed, "delete bucket failed")
}
return mq.ReplyOK(coormsg.NewDeleteBucketResp())
}

View File

@ -0,0 +1,100 @@
package services
import (
"database/sql"
"errors"
"github.com/jmoiron/sqlx"
"gitlink.org.cn/cloudream/common/consts/errorcode"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkg/mq"
mymq "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message"
coormsg "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message/coordinator"
)
func (svc *Service) PreUploadEcObject(msg *coormsg.PreUploadEcObject) (*coormsg.PreUploadEcResp, *mq.CodeMessage) {
// 判断同名对象是否存在。等到UploadRepObject时再判断一次。
// 此次的判断只作为参考具体是否成功还是看UploadRepObject的结果
isBucketAvai, err := svc.db.Bucket().IsAvailable(svc.db.SQLCtx(), msg.BucketID, msg.UserID)
if err != nil {
logger.WithField("BucketID", msg.BucketID).
Warnf("check bucket available failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "check bucket available failed")
}
if !isBucketAvai {
logger.WithField("BucketID", msg.BucketID).
Warnf("bucket is not available to user")
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "bucket is not available to user")
}
_, err = svc.db.Object().GetByName(svc.db.SQLCtx(), msg.BucketID, msg.ObjectName)
if err == nil {
logger.WithField("BucketID", msg.BucketID).
WithField("ObjectName", msg.ObjectName).
Warnf("object with given Name and BucketID already exists")
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "object with given Name and BucketID already exists")
}
if !errors.Is(err, sql.ErrNoRows) {
logger.WithField("BucketID", msg.BucketID).
WithField("ObjectName", msg.ObjectName).
Warnf("get object by name failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "get object by name failed")
}
//查询用户可用的节点IP
nodes, err := svc.db.Node().GetUserNodes(svc.db.SQLCtx(), msg.UserID)
if err != nil {
logger.WithField("UserID", msg.UserID).
Warnf("query user nodes failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "query user nodes failed")
}
// 查询客户端所属节点
foundBelongNode := true
belongNode, err := svc.db.Node().GetByExternalIP(svc.db.SQLCtx(), msg.ClientExternalIP)
if err == sql.ErrNoRows {
foundBelongNode = false
} else if err != nil {
logger.WithField("ClientExternalIP", msg.ClientExternalIP).
Warnf("query client belong node failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "query client belong node failed")
}
var respNodes []mymq.RespNode
for _, node := range nodes {
respNodes = append(respNodes, mymq.NewRespNode(
node.NodeID,
node.ExternalIP,
node.LocalIP,
// LocationID 相同则认为是在同一个地域
foundBelongNode && belongNode.LocationID == node.LocationID,
))
}
//查询纠删码参数
ec, err := svc.db.Ec().GetEc(svc.db.SQLCtx(), msg.EcName)
if err != nil {
logger.WithField("Ec", msg.EcName).
Warnf("check ec type failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadEcResp](errorcode.OperationFailed, "check bucket available failed")
}
ecc := mymq.NewEc(ec.EcID, ec.Name, ec.EcK, ec.EcN)
return mq.ReplyOK(coormsg.NewPreUploadEcResp(respNodes, ecc))
}
func (svc *Service) CreateEcObject(msg *coormsg.CreateEcObject) (*coormsg.CreateObjectResp, *mq.CodeMessage) {
var objID int64
err := svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
var err error
objID, err = svc.db.Object().CreateEcObject(tx, msg.BucketID, msg.ObjectName, msg.FileSize, msg.UserID, msg.NodeIDs, msg.Hashes, msg.EcName, msg.DirName)
return err
})
if err != nil {
logger.WithField("BucketName", msg.BucketID).
WithField("ObjectName", msg.ObjectName).
Warnf("create rep object failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.CreateObjectResp](errorcode.OperationFailed, "create rep object failed")
}
return mq.ReplyOK(coormsg.NewCreateObjectResp(objID))
}

358
internal/services/object.go Normal file
View File

@ -0,0 +1,358 @@
package services
import (
"database/sql"
"errors"
"github.com/jmoiron/sqlx"
"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/consts/errorcode"
"gitlink.org.cn/cloudream/common/models"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkg/mq"
"gitlink.org.cn/cloudream/storage-common/pkgs/db/model"
mymq "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message"
coormsg "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message/coordinator"
scevt "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message/scanner/event"
)
func (svc *Service) GetObjectsByDirName(msg *coormsg.GetObjectsByDirName) (*coormsg.GetObjectsResp, *mq.CodeMessage) {
//查询dirName下所有文件
objects, err := svc.db.Object().GetByDirName(svc.db.SQLCtx(), msg.DirName)
if err != nil {
logger.WithField("DirName", msg.DirName).
Warnf("query dirname failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.GetObjectsResp](errorcode.OperationFailed, "get objects failed")
}
return mq.ReplyOK(coormsg.NewGetObjectsResp(objects))
}
func (svc *Service) PreDownloadObject(msg *coormsg.PreDownloadObject) (*coormsg.PreDownloadObjectResp, *mq.CodeMessage) {
// 查询文件对象
object, err := svc.db.Object().GetUserObject(svc.db.SQLCtx(), msg.UserID, msg.ObjectID)
if err != nil {
logger.WithField("ObjectID", msg.ObjectID).
Warnf("query Object failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreDownloadObjectResp](errorcode.OperationFailed, "query Object failed")
}
// 查询客户端所属节点
foundBelongNode := true
belongNode, err := svc.db.Node().GetByExternalIP(svc.db.SQLCtx(), msg.ClientExternalIP)
if err == sql.ErrNoRows {
foundBelongNode = false
} else if err != nil {
logger.WithField("ClientExternalIP", msg.ClientExternalIP).
Warnf("query client belong node failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreDownloadObjectResp](errorcode.OperationFailed, "query client belong node failed")
}
logger.Debugf("client address %s is at location %d", msg.ClientExternalIP, belongNode.LocationID)
//-若redundancy是rep查询对象副本表, 获得FileHash
if object.Redundancy == models.RedundancyRep {
objectRep, err := svc.db.ObjectRep().GetByID(svc.db.SQLCtx(), object.ObjectID)
if err != nil {
logger.WithField("ObjectID", object.ObjectID).
Warnf("get ObjectRep failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreDownloadObjectResp](errorcode.OperationFailed, "query ObjectRep failed")
}
// 注由于采用了IPFS存储因此每个备份文件的FileHash都是一样的
nodes, err := svc.db.Cache().FindCachingFileUserNodes(svc.db.SQLCtx(), msg.UserID, objectRep.FileHash)
if err != nil {
logger.WithField("FileHash", objectRep.FileHash).
Warnf("query Cache failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreDownloadObjectResp](errorcode.OperationFailed, "query Cache failed")
}
var respNodes []mymq.RespNode
for _, node := range nodes {
respNodes = append(respNodes, mymq.NewRespNode(
node.NodeID,
node.ExternalIP,
node.LocalIP,
// LocationID 相同则认为是在同一个地域
foundBelongNode && belongNode.LocationID == node.LocationID,
))
}
return mq.ReplyOK(coormsg.NewPreDownloadObjectResp(
object.FileSize,
mymq.NewRespRepRedundancyData(objectRep.FileHash, respNodes),
))
} else {
// TODO 参考上面进行重写
ecName := object.Redundancy
blocks, err := svc.db.QueryObjectBlock(object.ObjectID)
if err != nil {
logger.WithField("ObjectID", object.ObjectID).
Warnf("query Blocks failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreDownloadObjectResp](errorcode.OperationFailed, "query Blocks failed")
}
logger.Debugf(blocks[4].BlockHash)
//查询纠删码参数
ec, err := svc.db.Ec().GetEc(svc.db.SQLCtx(), ecName)
ecc := mymq.NewEc(ec.EcID, ec.Name, ec.EcK, ec.EcN)
//查询每个编码块存放的所有节点
respNodes := make([][]mymq.RespNode, len(blocks))
for i := 0; i < len(blocks); i++ {
nodes, err := svc.db.Cache().FindCachingFileUserNodes(svc.db.SQLCtx(), msg.UserID, blocks[i].BlockHash)
if err != nil {
logger.WithField("FileHash", blocks[i].BlockHash).
Warnf("query Cache failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreDownloadObjectResp](errorcode.OperationFailed, "query Cache failed")
}
var nd []mymq.RespNode
for _, node := range nodes {
nd = append(nd, mymq.NewRespNode(
node.NodeID,
node.ExternalIP,
node.LocalIP,
// LocationID 相同则认为是在同一个地域
foundBelongNode && belongNode.LocationID == node.LocationID,
))
}
respNodes[i] = nd
logger.Debugf("##%d\n", i)
}
var blockss []mymq.RespObjectBlock
for i := 0; i < len(blocks); i++ {
blockss = append(blockss, mymq.NewRespObjectBlock(
blocks[i].InnerID,
blocks[i].BlockHash,
))
}
return mq.ReplyOK(coormsg.NewPreDownloadObjectResp(
object.FileSize,
mymq.NewRespEcRedundancyData(ecc, blockss, respNodes),
))
}
}
func (svc *Service) PreUploadRepObject(msg *coormsg.PreUploadRepObject) (*coormsg.PreUploadResp, *mq.CodeMessage) {
// 判断同名对象是否存在。等到UploadRepObject时再判断一次。
// 此次的判断只作为参考具体是否成功还是看UploadRepObject的结果
isBucketAvai, err := svc.db.Bucket().IsAvailable(svc.db.SQLCtx(), msg.BucketID, msg.UserID)
if err != nil {
logger.WithField("BucketID", msg.BucketID).
Warnf("check bucket available failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadResp](errorcode.OperationFailed, "check bucket available failed")
}
if !isBucketAvai {
logger.WithField("BucketID", msg.BucketID).
Warnf("bucket is not available to user")
return mq.ReplyFailed[coormsg.PreUploadResp](errorcode.OperationFailed, "bucket is not available to user")
}
_, err = svc.db.Object().GetByName(svc.db.SQLCtx(), msg.BucketID, msg.ObjectName)
if err == nil {
logger.WithField("BucketID", msg.BucketID).
WithField("ObjectName", msg.ObjectName).
Warnf("object with given Name and BucketID already exists")
return mq.ReplyFailed[coormsg.PreUploadResp](errorcode.OperationFailed, "object with given Name and BucketID already exists")
}
if !errors.Is(err, sql.ErrNoRows) {
logger.WithField("BucketID", msg.BucketID).
WithField("ObjectName", msg.ObjectName).
Warnf("get object by name failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadResp](errorcode.OperationFailed, "get object by name failed")
}
//查询用户可用的节点IP
nodes, err := svc.db.Node().GetUserNodes(svc.db.SQLCtx(), msg.UserID)
if err != nil {
logger.WithField("UserID", msg.UserID).
Warnf("query user nodes failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadResp](errorcode.OperationFailed, "query user nodes failed")
}
// 查询客户端所属节点
foundBelongNode := true
belongNode, err := svc.db.Node().GetByExternalIP(svc.db.SQLCtx(), msg.ClientExternalIP)
if err == sql.ErrNoRows {
foundBelongNode = false
} else if err != nil {
logger.WithField("ClientExternalIP", msg.ClientExternalIP).
Warnf("query client belong node failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUploadResp](errorcode.OperationFailed, "query client belong node failed")
}
var respNodes []mymq.RespNode
for _, node := range nodes {
respNodes = append(respNodes, mymq.NewRespNode(
node.NodeID,
node.ExternalIP,
node.LocalIP,
// LocationID 相同则认为是在同一个地域
foundBelongNode && belongNode.LocationID == node.LocationID,
))
}
return mq.ReplyOK(coormsg.NewPreUploadResp(respNodes))
}
func (svc *Service) CreateRepObject(msg *coormsg.CreateRepObject) (*coormsg.CreateObjectResp, *mq.CodeMessage) {
var objID int64
err := svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
var err error
objID, err = svc.db.Object().CreateRepObject(tx, msg.BucketID, msg.ObjectName, msg.FileSize, msg.RepCount, msg.NodeIDs, msg.FileHash, msg.DirName)
return err
})
if err != nil {
logger.WithField("BucketName", msg.BucketID).
WithField("ObjectName", msg.ObjectName).
Warnf("create rep object failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.CreateObjectResp](errorcode.OperationFailed, "create rep object failed")
}
// 紧急任务
err = svc.scanner.PostEvent(scevt.NewCheckRepCount([]string{msg.FileHash}), true, true)
if err != nil {
logger.Warnf("post event to scanner failed, but this will not affect creating, err: %s", err.Error())
}
return mq.ReplyOK(coormsg.NewCreateObjectResp(objID))
}
func (svc *Service) PreUpdateRepObject(msg *coormsg.PreUpdateRepObject) (*coormsg.PreUpdateRepObjectResp, *mq.CodeMessage) {
// TODO 检查用户是否有Object的权限
// 获取对象信息
obj, err := svc.db.Object().GetByID(svc.db.SQLCtx(), msg.ObjectID)
if err != nil {
logger.WithField("ObjectID", msg.ObjectID).
Warnf("get object failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUpdateRepObjectResp](errorcode.OperationFailed, "get object failed")
}
if obj.Redundancy != models.RedundancyRep {
logger.WithField("ObjectID", msg.ObjectID).
Warnf("this object is not a rep object")
return mq.ReplyFailed[coormsg.PreUpdateRepObjectResp](errorcode.OperationFailed, "this object is not a rep object")
}
// 获取对象Rep信息
objRep, err := svc.db.ObjectRep().GetByID(svc.db.SQLCtx(), msg.ObjectID)
if err != nil {
logger.WithField("ObjectID", msg.ObjectID).
Warnf("get object rep failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUpdateRepObjectResp](errorcode.OperationFailed, "get object rep failed")
}
//查询用户可用的节点IP
nodes, err := svc.db.Node().GetUserNodes(svc.db.SQLCtx(), msg.UserID)
if err != nil {
logger.WithField("UserID", msg.UserID).
Warnf("query user nodes failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUpdateRepObjectResp](errorcode.OperationFailed, "query user nodes failed")
}
// 查询客户端所属节点
foundBelongNode := true
belongNode, err := svc.db.Node().GetByExternalIP(svc.db.SQLCtx(), msg.ClientExternalIP)
if err == sql.ErrNoRows {
foundBelongNode = false
} else if err != nil {
logger.WithField("ClientExternalIP", msg.ClientExternalIP).
Warnf("query client belong node failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUpdateRepObjectResp](errorcode.OperationFailed, "query client belong node failed")
}
// 查询保存了旧文件的节点信息
cachingNodes, err := svc.db.Cache().FindCachingFileUserNodes(svc.db.SQLCtx(), msg.UserID, objRep.FileHash)
if err != nil {
logger.Warnf("find caching file user nodes failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreUpdateRepObjectResp](errorcode.OperationFailed, "find caching file user nodes failed")
}
var retNodes []coormsg.PreUpdateRepObjectRespNode
for _, node := range nodes {
retNodes = append(retNodes, coormsg.NewPreUpdateRepObjectRespNode(
node.NodeID,
node.ExternalIP,
node.LocalIP,
// LocationID 相同则认为是在同一个地域
foundBelongNode && belongNode.LocationID == node.LocationID,
// 此节点存储了对象旧文件
lo.ContainsBy(cachingNodes, func(n model.Node) bool { return n.NodeID == node.NodeID }),
))
}
return mq.ReplyOK(coormsg.NewPreUpdateRepObjectResp(retNodes))
}
func (svc *Service) UpdateRepObject(msg *coormsg.UpdateRepObject) (*coormsg.UpdateRepObjectResp, *mq.CodeMessage) {
err := svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
return svc.db.Object().UpdateRepObject(tx, msg.ObjectID, msg.FileSize, msg.NodeIDs, msg.FileHash)
})
if err != nil {
logger.WithField("ObjectID", msg.ObjectID).
Warnf("update rep object failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.UpdateRepObjectResp](errorcode.OperationFailed, "update rep object failed")
}
// 紧急任务
err = svc.scanner.PostEvent(scevt.NewCheckRepCount([]string{msg.FileHash}), true, true)
if err != nil {
logger.Warnf("post event to scanner failed, but this will not affect updating, err: %s", err.Error())
}
return mq.ReplyOK(coormsg.NewUpdateRepObjectResp())
}
func (svc *Service) DeleteObject(msg *coormsg.DeleteObject) (*coormsg.DeleteObjectResp, *mq.CodeMessage) {
isAva, err := svc.db.Object().IsAvailable(svc.db.SQLCtx(), msg.UserID, msg.ObjectID)
if err != nil {
logger.WithField("UserID", msg.UserID).
WithField("ObjectID", msg.ObjectID).
Warnf("check object available failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.DeleteObjectResp](errorcode.OperationFailed, "check object available failed")
}
if !isAva {
logger.WithField("UserID", msg.UserID).
WithField("ObjectID", msg.ObjectID).
Warnf("object is not available to the user")
return mq.ReplyFailed[coormsg.DeleteObjectResp](errorcode.OperationFailed, "object is not available to the user")
}
err = svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
return svc.db.Object().SoftDelete(tx, msg.ObjectID)
})
if err != nil {
logger.WithField("UserID", msg.UserID).
WithField("ObjectID", msg.ObjectID).
Warnf("set object deleted failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.DeleteObjectResp](errorcode.OperationFailed, "set object deleted failed")
}
stgs, err := svc.db.StorageObject().FindObjectStorages(svc.db.SQLCtx(), msg.ObjectID)
if err != nil {
logger.Warnf("find object storages failed, but this will not affect the deleting, err: %s", err.Error())
return mq.ReplyOK(coormsg.NewDeleteObjectResp())
}
// 不追求及时、准确
if len(stgs) == 0 {
// 如果没有被引用直接投递CheckObject的任务
err := svc.scanner.PostEvent(scevt.NewCheckObject([]int64{msg.ObjectID}), false, false)
if err != nil {
logger.Warnf("post event to scanner failed, but this will not affect deleting, err: %s", err.Error())
}
logger.Debugf("post check object event")
} else {
// 有引用则让Agent去检查StorageObject
for _, stg := range stgs {
err := svc.scanner.PostEvent(scevt.NewAgentCheckStorage(stg.StorageID, []int64{msg.ObjectID}), false, false)
if err != nil {
logger.Warnf("post event to scanner failed, but this will not affect deleting, err: %s", err.Error())
}
}
logger.Debugf("post agent check storage event")
}
return mq.ReplyOK(coormsg.NewDeleteObjectResp())
}

View File

@ -0,0 +1,18 @@
package services
import (
mydb "gitlink.org.cn/cloudream/storage-common/pkgs/db"
sccli "gitlink.org.cn/cloudream/storage-common/pkgs/mq/client/scanner"
)
type Service struct {
db *mydb.DB
scanner *sccli.Client
}
func NewService(db *mydb.DB, scanner *sccli.Client) *Service {
return &Service{
db: db,
scanner: scanner,
}
}

View File

@ -0,0 +1,103 @@
package services
import (
"database/sql"
"github.com/jmoiron/sqlx"
"gitlink.org.cn/cloudream/common/consts/errorcode"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/storage-common/models"
"gitlink.org.cn/cloudream/common/pkg/mq"
coormsg "gitlink.org.cn/cloudream/storage-common/pkgs/mq/message/coordinator"
)
func (svc *Service) GetStorageInfo(msg *coormsg.GetStorageInfo) (*coormsg.GetStorageInfoResp, *mq.CodeMessage) {
stg, err := svc.db.Storage().GetUserStorage(svc.db.SQLCtx(), msg.UserID, msg.StorageID)
if err != nil {
logger.Warnf("getting user storage: %s", err.Error())
return nil, mq.Failed(errorcode.OperationFailed, "get user storage failed")
}
return mq.ReplyOK(coormsg.NewGetStorageInfoResp(stg.StorageID, stg.Name, stg.NodeID, stg.Directory, stg.State))
}
func (svc *Service) PreMoveObjectToStorage(msg *coormsg.PreMoveObjectToStorage) (*coormsg.PreMoveObjectToStorageResp, *mq.CodeMessage) {
// 查询用户关联的存储服务
stg, err := svc.db.Storage().GetUserStorage(svc.db.SQLCtx(), msg.UserID, msg.StorageID)
if err != nil {
logger.WithField("UserID", msg.UserID).
WithField("StorageID", msg.StorageID).
Warnf("get user Storage failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OperationFailed, "get user Storage failed")
}
// 查询文件对象
object, err := svc.db.Object().GetUserObject(svc.db.SQLCtx(), msg.UserID, msg.ObjectID)
if err != nil {
logger.WithField("ObjectID", msg.ObjectID).
Warnf("get user Object failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OperationFailed, "get user Object failed")
}
//-若redundancy是rep查询对象副本表, 获得FileHash
if object.Redundancy == models.RedundancyRep {
objectRep, err := svc.db.ObjectRep().GetByID(svc.db.SQLCtx(), object.ObjectID)
if err != nil {
logger.Warnf("get ObjectRep failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OperationFailed, "get ObjectRep failed")
}
return mq.ReplyOK(coormsg.NewPreMoveObjectToStorageRespBody(
stg.NodeID,
stg.Directory,
object,
models.NewRedundancyRepData(objectRep.FileHash),
))
} else {
// TODO 以EC_开头的Redundancy才是EC策略
ecName := object.Redundancy
blocks, err := svc.db.QueryObjectBlock(object.ObjectID)
if err != nil {
logger.WithField("ObjectID", object.ObjectID).
Warnf("query Blocks failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OperationFailed, "query Blocks failed")
}
//查询纠删码参数
ec, err := svc.db.Ec().GetEc(svc.db.SQLCtx(), ecName)
// TODO zkx 异常处理
ecc := models.NewEc(ec.EcID, ec.Name, ec.EcK, ec.EcN)
blockss := make([]models.ObjectBlock, len(blocks))
for i := 0; i < len(blocks); i++ {
blockss[i] = models.NewObjectBlock(
blocks[i].InnerID,
blocks[i].BlockHash,
)
}
return mq.ReplyOK(coormsg.NewPreMoveObjectToStorageRespBody(
stg.NodeID,
stg.Directory,
object,
models.NewRedundancyEcData(ecc, blockss),
))
}
}
func (svc *Service) MoveObjectToStorage(msg *coormsg.MoveObjectToStorage) (*coormsg.MoveObjectToStorageResp, *mq.CodeMessage) {
// TODO: 对于的storage中已经存在的文件直接覆盖已有文件
err := svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
return svc.db.StorageObject().MoveObjectTo(tx, msg.ObjectID, msg.StorageID, msg.UserID)
})
if err != nil {
logger.WithField("UserID", msg.UserID).
WithField("ObjectID", msg.ObjectID).
WithField("StorageID", msg.StorageID).
Warnf("user move object to storage failed, err: %s", err.Error())
return mq.ReplyFailed[coormsg.MoveObjectToStorageResp](errorcode.OperationFailed, "user move object to storage failed")
}
return mq.ReplyOK(coormsg.NewMoveObjectToStorageResp())
}

20
magefiles/magefile.go Normal file
View File

@ -0,0 +1,20 @@
//go:build mage
package main
import (
"magefiles"
//mage:import
_ "magefiles/targets"
)
var Default = Build
func Build() error {
return magefiles.Build(magefiles.BuildArgs{
OutputName: "coordinator",
OutputDir: "../../build/coordinator",
AssetsDir: "assets",
})
}

64
main.go Normal file
View File

@ -0,0 +1,64 @@
package main
import (
"fmt"
"os"
"gitlink.org.cn/cloudream/common/pkg/logger"
log "gitlink.org.cn/cloudream/common/pkg/logger"
mydb "gitlink.org.cn/cloudream/storage-common/pkgs/db"
sccli "gitlink.org.cn/cloudream/storage-common/pkgs/mq/client/scanner"
rasvr "gitlink.org.cn/cloudream/storage-common/pkgs/mq/server/coordinator"
"gitlink.org.cn/cloudream/storage-coordinator/internal/config"
"gitlink.org.cn/cloudream/storage-coordinator/internal/services"
)
func main() {
err := config.Init()
if err != nil {
fmt.Printf("init config failed, err: %s", err.Error())
os.Exit(1)
}
err = logger.Init(&config.Cfg().Logger)
if err != nil {
fmt.Printf("init logger failed, err: %s", err.Error())
os.Exit(1)
}
db, err := mydb.NewDB(&config.Cfg().DB)
if err != nil {
log.Fatalf("new db failed, err: %s", err.Error())
}
scanner, err := sccli.NewClient(&config.Cfg().RabbitMQ)
if err != nil {
log.Fatalf("new scanner client failed, err: %s", err.Error())
}
coorSvr, err := rasvr.NewServer(services.NewService(db, scanner), &config.Cfg().RabbitMQ)
if err != nil {
log.Fatalf("new coordinator server failed, err: %s", err.Error())
}
coorSvr.OnError = func(err error) {
log.Warnf("coordinator server err: %s", err.Error())
}
// 启动服务
go serveCoorServer(coorSvr)
forever := make(chan bool)
<-forever
}
func serveCoorServer(server *rasvr.Server) {
log.Info("start serving command server")
err := server.Serve()
if err != nil {
log.Errorf("command server stopped with error: %s", err.Error())
}
log.Info("command server stopped")
}