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 12 kB

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

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