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 3.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
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101
  1. package services
  2. import (
  3. log "github.com/sirupsen/logrus"
  4. ramsg "gitlink.org.cn/cloudream/rabbitmq/message"
  5. coormsg "gitlink.org.cn/cloudream/rabbitmq/message/coordinator"
  6. "gitlink.org.cn/cloudream/utils"
  7. "gitlink.org.cn/cloudream/utils/consts"
  8. "gitlink.org.cn/cloudream/utils/consts/errorcode"
  9. )
  10. func (service *Service) Move(msg *coormsg.MoveObjectToStorage) *coormsg.MoveObjectToStorageResp {
  11. //查询数据库,获取冗余类型,冗余参数
  12. //jh:使用command中的bucketname和objectname查询对象表,获得redundancy,EcName,fileSizeInBytes
  13. //-若redundancy是rep,查询对象副本表, 获得repHash
  14. //--ids :={0}
  15. //--hashs := {repHash}
  16. //-若redundancy是ec,查询对象编码块表,获得blockHashs, ids(innerID),
  17. //--查询缓存表,获得每个hash的nodeIps、TempOrPins、Times
  18. //--查询节点延迟表,得到command.destination与各个nodeIps的的延迟,存到一个map类型中(Delay)
  19. //--kx:根据查出来的hash/hashs、nodeIps、TempOrPins、Times(移动/读取策略)、Delay确定hashs、ids
  20. // TODO 需要在StorageData中增加记录
  21. // 查询用户关联的存储服务
  22. stg, err := service.db.QueryUserStorage(msg.Body.UserID, msg.Body.StorageID)
  23. if err != nil {
  24. log.WithField("UserID", msg.Body.UserID).
  25. WithField("StorageID", msg.Body.StorageID).
  26. Warnf("query storage directory failed, err: %s", err.Error())
  27. return ramsg.ReplyFailed[coormsg.MoveObjectToStorageResp](errorcode.OPERATION_FAILED, "query storage directory failed")
  28. }
  29. // 查询文件对象
  30. object, err := service.db.QueryObjectByID(msg.Body.ObjectID)
  31. if err != nil {
  32. log.WithField("ObjectID", msg.Body.ObjectID).
  33. Warnf("query Object failed, err: %s", err.Error())
  34. return ramsg.ReplyFailed[coormsg.MoveObjectToStorageResp](errorcode.OPERATION_FAILED, "query Object failed")
  35. }
  36. //-若redundancy是rep,查询对象副本表, 获得repHash
  37. var hashs []string
  38. ids := []int{0}
  39. if object.Redundancy == consts.REDUNDANCY_REP {
  40. objectRep, err := service.db.QueryObjectRep(object.ObjectID)
  41. if err != nil {
  42. log.Warnf("query ObjectRep failed, err: %s", err.Error())
  43. return ramsg.ReplyFailed[coormsg.MoveObjectToStorageResp](errorcode.OPERATION_FAILED, "query ObjectRep failed")
  44. }
  45. hashs = append(hashs, objectRep.RepHash)
  46. } else {
  47. blockHashs, err := service.db.QueryObjectBlock(object.ObjectID)
  48. if err != nil {
  49. log.Warnf("query ObjectBlock failed, err: %s", err.Error())
  50. return ramsg.ReplyFailed[coormsg.MoveObjectToStorageResp](errorcode.OPERATION_FAILED, "query ObjectBlock failed")
  51. }
  52. ecPolicies := *utils.GetEcPolicy()
  53. ecPolicy := ecPolicies[*object.ECName]
  54. ecN := ecPolicy.GetN()
  55. ecK := ecPolicy.GetK()
  56. ids = make([]int, ecK)
  57. for i := 0; i < ecN; i++ {
  58. hashs = append(hashs, "-1")
  59. }
  60. for i := 0; i < ecK; i++ {
  61. ids[i] = i
  62. }
  63. hashs = make([]string, ecN)
  64. for _, tt := range blockHashs {
  65. id := tt.InnerID
  66. hash := tt.BlockHash
  67. hashs[id] = hash
  68. }
  69. //--查询缓存表,获得每个hash的nodeIps、TempOrPins、Times
  70. /*for id,hash := range blockHashs{
  71. //type Cache struct {NodeIP string,TempOrPin bool,Cachetime string}
  72. Cache := Query_Cache(hash)
  73. //利用Time_trans()函数可将Cache[i].Cachetime转化为时间戳格式
  74. //--查询节点延迟表,得到command.Destination与各个nodeIps的延迟,存到一个map类型中(Delay)
  75. Delay := make(map[string]int) // 延迟集合
  76. for i:=0; i<len(Cache); i++{
  77. Delay[Cache[i].NodeIP] = Query_NodeDelay(Destination, Cache[i].NodeIP)
  78. }
  79. //--kx:根据查出来的hash/hashs、nodeIps、TempOrPins、Times(移动/读取策略)、Delay确定hashs、ids
  80. }*/
  81. }
  82. return ramsg.ReplyOK(coormsg.NewMoveObjectToStorageRespBody(
  83. stg.NodeID,
  84. stg.Directory,
  85. object.Redundancy,
  86. object.ECName,
  87. hashs,
  88. ids,
  89. object.FileSizeInBytes,
  90. ))
  91. }

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