You can not select more than 25 topics Topics must start with a chinese character,a letter or number, can include dashes ('-') and can be up to 35 characters long.

storage.go 4.8 kB

2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
2 years ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116
  1. package services
  2. import (
  3. "database/sql"
  4. "github.com/jmoiron/sqlx"
  5. "gitlink.org.cn/cloudream/common/consts"
  6. "gitlink.org.cn/cloudream/common/consts/errorcode"
  7. log "gitlink.org.cn/cloudream/common/pkg/logger"
  8. "gitlink.org.cn/cloudream/common/utils"
  9. ramsg "gitlink.org.cn/cloudream/rabbitmq/message"
  10. coormsg "gitlink.org.cn/cloudream/rabbitmq/message/coordinator"
  11. )
  12. func (svc *Service) PreMoveObjectToStorage(msg *coormsg.PreMoveObjectToStorage) *coormsg.PreMoveObjectToStorageResp {
  13. //查询数据库,获取冗余类型,冗余参数
  14. //jh:使用command中的bucketname和objectname查询对象表,获得redundancy,EcName,fileSize
  15. //-若redundancy是rep,查询对象副本表, 获得repHash
  16. //--ids :={0}
  17. //--hashs := {repHash}
  18. //-若redundancy是ec,查询对象编码块表,获得blockHashs, ids(innerID),
  19. //--查询缓存表,获得每个hash的nodeIps、TempOrPins、Times
  20. //--查询节点延迟表,得到command.destination与各个nodeIps的的延迟,存到一个map类型中(Delay)
  21. //--kx:根据查出来的hash/hashs、nodeIps、TempOrPins、Times(移动/读取策略)、Delay确定hashs、ids
  22. // 查询用户关联的存储服务
  23. stg, err := svc.db.Storage().GetUserStorage(svc.db.SQLCtx(), msg.Body.UserID, msg.Body.StorageID)
  24. if err != nil {
  25. log.WithField("UserID", msg.Body.UserID).
  26. WithField("StorageID", msg.Body.StorageID).
  27. Warnf("get user Storage failed, err: %s", err.Error())
  28. return ramsg.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OPERATION_FAILED, "get user Storage failed")
  29. }
  30. // 查询文件对象
  31. object, err := svc.db.Object().GetUserObject(svc.db.SQLCtx(), msg.Body.UserID, msg.Body.ObjectID)
  32. if err != nil {
  33. log.WithField("ObjectID", msg.Body.ObjectID).
  34. Warnf("get user Object failed, err: %s", err.Error())
  35. return ramsg.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OPERATION_FAILED, "get user Object failed")
  36. }
  37. //-若redundancy是rep,查询对象副本表, 获得FileHash
  38. if object.Redundancy == consts.REDUNDANCY_REP {
  39. objectRep, err := svc.db.ObjectRep().GetByID(svc.db.SQLCtx(), object.ObjectID)
  40. if err != nil {
  41. log.Warnf("get ObjectRep failed, err: %s", err.Error())
  42. return ramsg.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OPERATION_FAILED, "get ObjectRep failed")
  43. }
  44. return ramsg.ReplyOK(coormsg.NewPreMoveObjectToStorageRespBody(
  45. stg.NodeID,
  46. stg.Directory,
  47. object.FileSize,
  48. object.Redundancy,
  49. ramsg.NewObjectRepInfo(objectRep.FileHash),
  50. ))
  51. } else {
  52. // TODO 以EC_开头的Redundancy才是EC策略
  53. var hashs []string
  54. ids := []int{0}
  55. blockHashs, err := svc.db.QueryObjectBlock(object.ObjectID)
  56. if err != nil {
  57. log.Warnf("query ObjectBlock failed, err: %s", err.Error())
  58. return ramsg.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OPERATION_FAILED, "query ObjectBlock failed")
  59. }
  60. ecPolicies := *utils.GetEcPolicy()
  61. ecPolicy := ecPolicies[object.Redundancy]
  62. ecN := ecPolicy.GetN()
  63. ecK := ecPolicy.GetK()
  64. ids = make([]int, ecK)
  65. for i := 0; i < ecN; i++ {
  66. hashs = append(hashs, "-1")
  67. }
  68. for i := 0; i < ecK; i++ {
  69. ids[i] = i
  70. }
  71. hashs = make([]string, ecN)
  72. for _, tt := range blockHashs {
  73. id := tt.InnerID
  74. hash := tt.BlockHash
  75. hashs[id] = hash
  76. }
  77. //--查询缓存表,获得每个hash的nodeIps、TempOrPins、Times
  78. /*for id,hash := range blockHashs{
  79. //type Cache struct {NodeIP string,TempOrPin bool,Cachetime string}
  80. Cache := Query_Cache(hash)
  81. //利用Time_trans()函数可将Cache[i].Cachetime转化为时间戳格式
  82. //--查询节点延迟表,得到command.Destination与各个nodeIps的延迟,存到一个map类型中(Delay)
  83. Delay := make(map[string]int) // 延迟集合
  84. for i:=0; i<len(Cache); i++{
  85. Delay[Cache[i].NodeIP] = Query_NodeDelay(Destination, Cache[i].NodeIP)
  86. }
  87. //--kx:根据查出来的hash/hashs、nodeIps、TempOrPins、Times(移动/读取策略)、Delay确定hashs、ids
  88. }*/
  89. return ramsg.ReplyFailed[coormsg.PreMoveObjectToStorageResp](errorcode.OPERATION_FAILED, "not implement yet!")
  90. }
  91. }
  92. func (svc *Service) MoveObjectToStorage(msg *coormsg.MoveObjectToStorage) *coormsg.MoveObjectToStorageResp {
  93. err := svc.db.DoTx(sql.LevelDefault, func(tx *sqlx.Tx) error {
  94. return svc.db.StorageObject().MoveObjectTo(tx, msg.Body.UserID, msg.Body.ObjectID, msg.Body.StorageID)
  95. })
  96. if err != nil {
  97. log.WithField("UserID", msg.Body.UserID).
  98. WithField("ObjectID", msg.Body.ObjectID).
  99. WithField("StorageID", msg.Body.StorageID).
  100. Warnf("user move object to storage failed, err: %s", err.Error())
  101. return ramsg.ReplyFailed[coormsg.MoveObjectToStorageResp](errorcode.OPERATION_FAILED, "user move object to storage failed")
  102. }
  103. return ramsg.ReplyOK(coormsg.NewMoveObjectToStorageRespBody())
  104. }

本项目旨在将云际存储公共基础设施化,使个人及企业可低门槛使用高效的云际存储服务(安装开箱即用云际存储客户端即可,无需关注其他组件的部署),同时支持用户灵活便捷定制云际存储的功能细节。