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.

object.go 11 kB

1 year ago
1 year ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393
  1. package db2
  2. import (
  3. "fmt"
  4. "time"
  5. "gorm.io/gorm/clause"
  6. cdssdk "gitlink.org.cn/cloudream/common/sdks/storage"
  7. stgmod "gitlink.org.cn/cloudream/storage/common/models"
  8. "gitlink.org.cn/cloudream/storage/common/pkgs/db2/model"
  9. coormq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/coordinator"
  10. )
  11. type ObjectDB struct {
  12. *DB
  13. }
  14. func (db *DB) Object() *ObjectDB {
  15. return &ObjectDB{DB: db}
  16. }
  17. func (db *ObjectDB) GetByID(ctx SQLContext, objectID cdssdk.ObjectID) (model.Object, error) {
  18. var ret cdssdk.Object
  19. err := ctx.Table("Object").Where("ObjectID = ?", objectID).First(&ret).Error
  20. return ret, err
  21. }
  22. func (db *ObjectDB) GetWithPathPrefix(ctx SQLContext, packageID cdssdk.PackageID, pathPrefix string) ([]cdssdk.Object, error) {
  23. var ret []cdssdk.Object
  24. err := ctx.Table("Object").Where("PackageID = ? AND Path LIKE ?", packageID, pathPrefix+"%").Order("ObjectID ASC").Find(&ret).Error
  25. return ret, err
  26. }
  27. func (db *ObjectDB) BatchTestObjectID(ctx SQLContext, objectIDs []cdssdk.ObjectID) (map[cdssdk.ObjectID]bool, error) {
  28. if len(objectIDs) == 0 {
  29. return make(map[cdssdk.ObjectID]bool), nil
  30. }
  31. var avaiIDs []cdssdk.ObjectID
  32. err := ctx.Table("Object").Where("ObjectID IN ?", objectIDs).Pluck("ObjectID", &avaiIDs).Error
  33. if err != nil {
  34. return nil, err
  35. }
  36. avaiIDMap := make(map[cdssdk.ObjectID]bool)
  37. for _, pkgID := range avaiIDs {
  38. avaiIDMap[pkgID] = true
  39. }
  40. return avaiIDMap, nil
  41. }
  42. func (db *ObjectDB) BatchGet(ctx SQLContext, objectIDs []cdssdk.ObjectID) ([]model.Object, error) {
  43. if len(objectIDs) == 0 {
  44. return nil, nil
  45. }
  46. var objs []cdssdk.Object
  47. err := ctx.Table("Object").Where("ObjectID IN ?", objectIDs).Order("ObjectID ASC").Find(&objs).Error
  48. if err != nil {
  49. return nil, err
  50. }
  51. return objs, nil
  52. }
  53. func (db *ObjectDB) BatchGetByPackagePath(ctx SQLContext, pkgID cdssdk.PackageID, pathes []string) ([]cdssdk.Object, error) {
  54. if len(pathes) == 0 {
  55. return nil, nil
  56. }
  57. var objs []cdssdk.Object
  58. err := ctx.Table("Object").Where("PackageID = ? AND Path IN ?", pkgID, pathes).Find(&objs).Error
  59. if err != nil {
  60. return nil, err
  61. }
  62. return objs, nil
  63. }
  64. func (db *ObjectDB) Create(ctx SQLContext, obj cdssdk.Object) (cdssdk.ObjectID, error) {
  65. err := ctx.Table("Object").Create(&obj).Error
  66. if err != nil {
  67. return 0, fmt.Errorf("insert object failed, err: %w", err)
  68. }
  69. return obj.ObjectID, nil
  70. }
  71. // 批量创建对象,创建完成后会填充ObjectID。
  72. func (db *ObjectDB) BatchCreate(ctx SQLContext, objs *[]cdssdk.Object) error {
  73. if len(*objs) == 0 {
  74. return nil
  75. }
  76. return ctx.Table("Object").Create(objs).Error
  77. }
  78. // 批量更新对象所有属性,objs中的对象必须包含ObjectID
  79. func (db *ObjectDB) BatchUpdate(ctx SQLContext, objs []cdssdk.Object) error {
  80. if len(objs) == 0 {
  81. return nil
  82. }
  83. return ctx.Clauses(clause.OnConflict{
  84. Columns: []clause.Column{{Name: "ObjectID"}},
  85. UpdateAll: true,
  86. }).Create(objs).Error
  87. }
  88. // 批量更新对象指定属性,objs中的对象只需设置需要更新的属性即可,但:
  89. // 1. 必须包含ObjectID
  90. // 2. 日期类型属性不能设置为0值
  91. func (db *ObjectDB) BatchUpdateColumns(ctx SQLContext, objs []cdssdk.Object, columns []string) error {
  92. if len(objs) == 0 {
  93. return nil
  94. }
  95. return ctx.Clauses(clause.OnConflict{
  96. Columns: []clause.Column{{Name: "ObjectID"}},
  97. DoUpdates: clause.AssignmentColumns(columns),
  98. }).Create(objs).Error
  99. }
  100. func (db *ObjectDB) GetPackageObjects(ctx SQLContext, packageID cdssdk.PackageID) ([]model.Object, error) {
  101. var ret []cdssdk.Object
  102. err := ctx.Table("Object").Where("PackageID = ?", packageID).Order("ObjectID ASC").Find(&ret).Error
  103. return ret, err
  104. }
  105. func (db *ObjectDB) GetPackageObjectDetails(ctx SQLContext, packageID cdssdk.PackageID) ([]stgmod.ObjectDetail, error) {
  106. var objs []cdssdk.Object
  107. err := ctx.Table("Object").Where("PackageID = ?", packageID).Order("ObjectID ASC").Find(&objs).Error
  108. if err != nil {
  109. return nil, fmt.Errorf("getting objects: %w", err)
  110. }
  111. // 获取所有的 ObjectBlock
  112. var allBlocks []stgmod.ObjectBlock
  113. err = ctx.Table("ObjectBlock").
  114. Select("ObjectBlock.*").
  115. Joins("JOIN Object ON ObjectBlock.ObjectID = Object.ObjectID").
  116. Where("Object.PackageID = ?", packageID).
  117. Order("ObjectBlock.ObjectID, `Index` ASC").
  118. Find(&allBlocks).Error
  119. if err != nil {
  120. return nil, fmt.Errorf("getting all object blocks: %w", err)
  121. }
  122. // 获取所有的 PinnedObject
  123. var allPinnedObjs []cdssdk.PinnedObject
  124. err = ctx.Table("PinnedObject").
  125. Select("PinnedObject.*").
  126. Joins("JOIN Object ON PinnedObject.ObjectID = Object.ObjectID").
  127. Where("Object.PackageID = ?", packageID).
  128. Order("PinnedObject.ObjectID").
  129. Find(&allPinnedObjs).Error
  130. if err != nil {
  131. return nil, fmt.Errorf("getting all pinned objects: %w", err)
  132. }
  133. details := make([]stgmod.ObjectDetail, len(objs))
  134. for i, obj := range objs {
  135. details[i] = stgmod.ObjectDetail{
  136. Object: obj,
  137. }
  138. }
  139. stgmod.DetailsFillObjectBlocks(details, allBlocks)
  140. stgmod.DetailsFillPinnedAt(details, allPinnedObjs)
  141. return details, nil
  142. }
  143. func (db *ObjectDB) GetObjectsIfAnyBlockOnStorage(ctx SQLContext, stgID cdssdk.StorageID) ([]cdssdk.Object, error) {
  144. var objs []cdssdk.Object
  145. err := ctx.Table("Object").Where("ObjectID IN (SELECT ObjectID FROM ObjectBlock WHERE StorageID = ?)", stgID).Order("ObjectID ASC").Find(&objs).Error
  146. if err != nil {
  147. return nil, fmt.Errorf("getting objects: %w", err)
  148. }
  149. return objs, nil
  150. }
  151. func (db *ObjectDB) BatchAdd(ctx SQLContext, packageID cdssdk.PackageID, adds []coormq.AddObjectEntry) ([]cdssdk.Object, error) {
  152. if len(adds) == 0 {
  153. return nil, nil
  154. }
  155. // 收集所有路径
  156. pathes := make([]string, 0, len(adds))
  157. for _, add := range adds {
  158. pathes = append(pathes, add.Path)
  159. }
  160. // 先查询要更新的对象,不存在也没关系
  161. existsObjs, err := db.BatchGetByPackagePath(ctx, packageID, pathes)
  162. if err != nil {
  163. return nil, fmt.Errorf("batch get object by path: %w", err)
  164. }
  165. existsObjsMap := make(map[string]cdssdk.Object)
  166. for _, obj := range existsObjs {
  167. existsObjsMap[obj.Path] = obj
  168. }
  169. var updatingObjs []cdssdk.Object
  170. var addingObjs []cdssdk.Object
  171. for i := range adds {
  172. o := cdssdk.Object{
  173. PackageID: packageID,
  174. Path: adds[i].Path,
  175. Size: adds[i].Size,
  176. FileHash: adds[i].FileHash,
  177. Redundancy: cdssdk.NewNoneRedundancy(), // 首次上传默认使用不分块的none模式
  178. CreateTime: adds[i].UploadTime,
  179. UpdateTime: adds[i].UploadTime,
  180. }
  181. e, ok := existsObjsMap[adds[i].Path]
  182. if ok {
  183. o.ObjectID = e.ObjectID
  184. o.CreateTime = e.CreateTime
  185. updatingObjs = append(updatingObjs, o)
  186. } else {
  187. addingObjs = append(addingObjs, o)
  188. }
  189. }
  190. // 先进行更新
  191. err = db.BatchUpdate(ctx, updatingObjs)
  192. if err != nil {
  193. return nil, fmt.Errorf("batch update objects: %w", err)
  194. }
  195. // 再执行插入,Create函数插入后会填充ObjectID
  196. err = db.BatchCreate(ctx, &addingObjs)
  197. if err != nil {
  198. return nil, fmt.Errorf("batch create objects: %w", err)
  199. }
  200. // 按照add参数的顺序返回结果
  201. affectedObjsMp := make(map[string]cdssdk.Object)
  202. for _, o := range updatingObjs {
  203. affectedObjsMp[o.Path] = o
  204. }
  205. for _, o := range addingObjs {
  206. affectedObjsMp[o.Path] = o
  207. }
  208. affectedObjs := make([]cdssdk.Object, 0, len(affectedObjsMp))
  209. affectedObjIDs := make([]cdssdk.ObjectID, 0, len(affectedObjsMp))
  210. for i := range adds {
  211. obj := affectedObjsMp[adds[i].Path]
  212. affectedObjs = append(affectedObjs, obj)
  213. affectedObjIDs = append(affectedObjIDs, obj.ObjectID)
  214. }
  215. if len(affectedObjIDs) > 0 {
  216. // 批量删除 ObjectBlock
  217. if err := ctx.Table("ObjectBlock").Where("ObjectID IN ?", affectedObjIDs).Delete(&stgmod.ObjectBlock{}).Error; err != nil {
  218. return nil, fmt.Errorf("batch delete object blocks: %w", err)
  219. }
  220. // 批量删除 PinnedObject
  221. if err := ctx.Table("PinnedObject").Where("ObjectID IN ?", affectedObjIDs).Delete(&cdssdk.PinnedObject{}).Error; err != nil {
  222. return nil, fmt.Errorf("batch delete pinned objects: %w", err)
  223. }
  224. }
  225. // 创建 ObjectBlock
  226. objBlocks := make([]stgmod.ObjectBlock, 0, len(adds))
  227. for i, add := range adds {
  228. for _, stgID := range add.StorageIDs {
  229. objBlocks = append(objBlocks, stgmod.ObjectBlock{
  230. ObjectID: affectedObjIDs[i],
  231. Index: 0,
  232. StorageID: stgID,
  233. FileHash: add.FileHash,
  234. })
  235. }
  236. }
  237. if err := db.ObjectBlock().BatchCreate(ctx, objBlocks); err != nil {
  238. return nil, fmt.Errorf("batch create object blocks: %w", err)
  239. }
  240. // 创建 Cache
  241. caches := make([]model.Cache, 0, len(adds))
  242. for _, add := range adds {
  243. for _, stgID := range add.StorageIDs {
  244. caches = append(caches, model.Cache{
  245. FileHash: add.FileHash,
  246. StorageID: stgID,
  247. CreateTime: time.Now(),
  248. Priority: 0,
  249. })
  250. }
  251. }
  252. if err := db.Cache().BatchCreate(ctx, caches); err != nil {
  253. return nil, fmt.Errorf("batch create caches: %w", err)
  254. }
  255. return affectedObjs, nil
  256. }
  257. func (db *ObjectDB) BatchUpdateRedundancy(ctx SQLContext, objs []coormq.UpdatingObjectRedundancy) error {
  258. if len(objs) == 0 {
  259. return nil
  260. }
  261. nowTime := time.Now()
  262. objIDs := make([]cdssdk.ObjectID, 0, len(objs))
  263. dummyObjs := make([]cdssdk.Object, 0, len(objs))
  264. for _, obj := range objs {
  265. objIDs = append(objIDs, obj.ObjectID)
  266. dummyObjs = append(dummyObjs, cdssdk.Object{
  267. ObjectID: obj.ObjectID,
  268. Redundancy: obj.Redundancy,
  269. CreateTime: nowTime, // 实际不会更新,只因为不能是0值
  270. UpdateTime: nowTime,
  271. })
  272. }
  273. err := db.Object().BatchUpdateColumns(ctx, dummyObjs, []string{"Redundancy", "UpdateTime"})
  274. if err != nil {
  275. return fmt.Errorf("batch update object redundancy: %w", err)
  276. }
  277. // 删除原本所有的编码块记录,重新添加
  278. err = db.ObjectBlock().BatchDeleteByObjectID(ctx, objIDs)
  279. if err != nil {
  280. return fmt.Errorf("batch delete object blocks: %w", err)
  281. }
  282. // 删除原本Pin住的Object。暂不考虑FileHash没有变化的情况
  283. err = db.PinnedObject().BatchDeleteByObjectID(ctx, objIDs)
  284. if err != nil {
  285. return fmt.Errorf("batch delete pinned object: %w", err)
  286. }
  287. blocks := make([]stgmod.ObjectBlock, 0, len(objs))
  288. for _, obj := range objs {
  289. blocks = append(blocks, obj.Blocks...)
  290. }
  291. err = db.ObjectBlock().BatchCreate(ctx, blocks)
  292. if err != nil {
  293. return fmt.Errorf("batch create object blocks: %w", err)
  294. }
  295. caches := make([]model.Cache, 0, len(objs))
  296. for _, obj := range objs {
  297. for _, blk := range obj.Blocks {
  298. caches = append(caches, model.Cache{
  299. FileHash: blk.FileHash,
  300. StorageID: blk.StorageID,
  301. CreateTime: nowTime,
  302. Priority: 0,
  303. })
  304. }
  305. }
  306. err = db.Cache().BatchCreate(ctx, caches)
  307. if err != nil {
  308. return fmt.Errorf("batch create object caches: %w", err)
  309. }
  310. pinneds := make([]cdssdk.PinnedObject, 0, len(objs))
  311. for _, obj := range objs {
  312. for _, p := range obj.PinnedAt {
  313. pinneds = append(pinneds, cdssdk.PinnedObject{
  314. ObjectID: obj.ObjectID,
  315. StorageID: p,
  316. CreateTime: nowTime,
  317. })
  318. }
  319. }
  320. err = db.PinnedObject().BatchTryCreate(ctx, pinneds)
  321. if err != nil {
  322. return fmt.Errorf("batch create pinned objects: %w", err)
  323. }
  324. return nil
  325. }
  326. func (db *ObjectDB) BatchDelete(ctx SQLContext, ids []cdssdk.ObjectID) error {
  327. if len(ids) == 0 {
  328. return nil
  329. }
  330. return ctx.Table("Object").Where("ObjectID IN ?", ids).Delete(&cdssdk.Object{}).Error
  331. }
  332. func (db *ObjectDB) DeleteInPackage(ctx SQLContext, packageID cdssdk.PackageID) error {
  333. return ctx.Table("Object").Where("PackageID = ?", packageID).Delete(&cdssdk.Object{}).Error
  334. }

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