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

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511
  1. package db2
  2. import (
  3. "fmt"
  4. "strings"
  5. "time"
  6. "gorm.io/gorm"
  7. "gorm.io/gorm/clause"
  8. cdssdk "gitlink.org.cn/cloudream/common/sdks/storage"
  9. stgmod "gitlink.org.cn/cloudream/storage2/common/models"
  10. "gitlink.org.cn/cloudream/storage2/common/pkgs/db2/model"
  11. coormq "gitlink.org.cn/cloudream/storage2/common/pkgs/mq/coordinator"
  12. )
  13. type ObjectDB struct {
  14. *DB
  15. }
  16. func (db *DB) Object() *ObjectDB {
  17. return &ObjectDB{DB: db}
  18. }
  19. func (db *ObjectDB) GetByID(ctx SQLContext, objectID cdssdk.ObjectID) (cdssdk.Object, error) {
  20. var ret cdssdk.Object
  21. err := ctx.Table("Object").Where("ObjectID = ?", objectID).First(&ret).Error
  22. return ret, err
  23. }
  24. func (db *ObjectDB) GetByPath(ctx SQLContext, packageID cdssdk.PackageID, path string) (cdssdk.Object, error) {
  25. var ret cdssdk.Object
  26. err := ctx.Table("Object").Where("PackageID = ? AND Path = ?", packageID, path).First(&ret).Error
  27. return ret, err
  28. }
  29. func (db *ObjectDB) GetByFullPath(ctx SQLContext, bktName string, pkgName string, path string) (cdssdk.Object, error) {
  30. var ret cdssdk.Object
  31. err := ctx.Table("Object").
  32. Joins("join Package on Package.PackageID = Object.PackageID and Package.Name = ?", pkgName).
  33. Joins("join Bucket on Bucket.BucketID = Package.BucketID and Bucket.Name = ?", bktName).
  34. Where("Object.Path = ?", path).First(&ret).Error
  35. return ret, err
  36. }
  37. func (db *ObjectDB) GetWithPathPrefix(ctx SQLContext, packageID cdssdk.PackageID, pathPrefix string) ([]cdssdk.Object, error) {
  38. var ret []cdssdk.Object
  39. err := ctx.Table("Object").Where("PackageID = ? AND Path LIKE ?", packageID, escapeLike("", "%", pathPrefix)).Order("ObjectID ASC").Find(&ret).Error
  40. return ret, err
  41. }
  42. // 查询结果将按照Path升序,而不是ObjectID升序
  43. func (db *ObjectDB) GetWithPathPrefixPaged(ctx SQLContext, packageID cdssdk.PackageID, pathPrefix string, startPath string, limit int) ([]cdssdk.Object, error) {
  44. var ret []cdssdk.Object
  45. err := ctx.Table("Object").Where("PackageID = ? AND Path > ? AND Path LIKE ?", packageID, startPath, pathPrefix+"%").Order("Path ASC").Limit(limit).Find(&ret).Error
  46. return ret, err
  47. }
  48. func (db *ObjectDB) GetByPrefixGrouped(ctx SQLContext, packageID cdssdk.PackageID, pathPrefix string) (objs []cdssdk.Object, commonPrefixes []string, err error) {
  49. type ObjectOrDir struct {
  50. cdssdk.Object
  51. IsObject bool `gorm:"IsObject"`
  52. Prefix string `gorm:"Prefix"`
  53. }
  54. sepCnt := strings.Count(pathPrefix, cdssdk.ObjectPathSeparator) + 1
  55. prefixStatm := fmt.Sprintf("Substring_Index(Path, '%s', %d)", cdssdk.ObjectPathSeparator, sepCnt)
  56. grouping := ctx.Table("Object").
  57. Select(fmt.Sprintf("%s as Prefix, Max(ObjectID) as ObjectID, %s = Path as IsObject", prefixStatm, prefixStatm)).
  58. Where("PackageID = ?", packageID).
  59. Where("Path like ?", pathPrefix+"%").
  60. Group("Prefix, IsObject").
  61. Order("Prefix ASC")
  62. var ret []ObjectOrDir
  63. err = ctx.Table("Object").
  64. Select("Grouped.IsObject, Grouped.Prefix, Object.*").
  65. Joins("right join (?) as Grouped on Object.ObjectID = Grouped.ObjectID and Grouped.IsObject = 1", grouping).
  66. Find(&ret).Error
  67. if err != nil {
  68. return
  69. }
  70. for _, o := range ret {
  71. if o.IsObject {
  72. objs = append(objs, o.Object)
  73. } else {
  74. commonPrefixes = append(commonPrefixes, o.Prefix+cdssdk.ObjectPathSeparator)
  75. }
  76. }
  77. return
  78. }
  79. func (db *ObjectDB) GetByPrefixGroupedPaged(ctx SQLContext, packageID cdssdk.PackageID, pathPrefix string, startPath string, limit int) (objs []cdssdk.Object, commonPrefixes []string, nextStartPath string, err error) {
  80. type ObjectOrDir struct {
  81. cdssdk.Object
  82. IsObject bool `gorm:"IsObject"`
  83. Prefix string `gorm:"Prefix"`
  84. }
  85. sepCnt := strings.Count(pathPrefix, cdssdk.ObjectPathSeparator) + 1
  86. prefixStatm := fmt.Sprintf("Substring_Index(Path, '%s', %d)", cdssdk.ObjectPathSeparator, sepCnt)
  87. grouping := ctx.Table("Object").
  88. Select(fmt.Sprintf("%s as Prefix, Max(ObjectID) as ObjectID, %s = Path as IsObject", prefixStatm, prefixStatm)).
  89. Where("PackageID = ?", packageID).
  90. Where("Path like ?", pathPrefix+"%").
  91. Group("Prefix, IsObject").
  92. Having("Prefix > ?", startPath).
  93. Limit(limit).
  94. Order("Prefix ASC")
  95. var ret []ObjectOrDir
  96. err = ctx.Table("Object").
  97. Select("Grouped.IsObject, Grouped.Prefix, Object.*").
  98. Joins("right join (?) as Grouped on Object.ObjectID = Grouped.ObjectID and Grouped.IsObject = 1", grouping).
  99. Find(&ret).Error
  100. if err != nil {
  101. return
  102. }
  103. for _, o := range ret {
  104. if o.IsObject {
  105. objs = append(objs, o.Object)
  106. } else {
  107. commonPrefixes = append(commonPrefixes, o.Prefix+cdssdk.ObjectPathSeparator)
  108. }
  109. nextStartPath = o.Prefix
  110. }
  111. return
  112. }
  113. func (db *ObjectDB) HasObjectWithPrefix(ctx SQLContext, packageID cdssdk.PackageID, pathPrefix string) (bool, error) {
  114. var obj cdssdk.Object
  115. err := ctx.Table("Object").Where("PackageID = ? AND Path LIKE ?", packageID, escapeLike("", "%", pathPrefix)).First(&obj).Error
  116. if err == nil {
  117. return true, nil
  118. }
  119. if err == gorm.ErrRecordNotFound {
  120. return false, nil
  121. }
  122. return false, err
  123. }
  124. func (db *ObjectDB) BatchTestObjectID(ctx SQLContext, objectIDs []cdssdk.ObjectID) (map[cdssdk.ObjectID]bool, error) {
  125. if len(objectIDs) == 0 {
  126. return make(map[cdssdk.ObjectID]bool), nil
  127. }
  128. var avaiIDs []cdssdk.ObjectID
  129. err := ctx.Table("Object").Where("ObjectID IN ?", objectIDs).Pluck("ObjectID", &avaiIDs).Error
  130. if err != nil {
  131. return nil, err
  132. }
  133. avaiIDMap := make(map[cdssdk.ObjectID]bool)
  134. for _, pkgID := range avaiIDs {
  135. avaiIDMap[pkgID] = true
  136. }
  137. return avaiIDMap, nil
  138. }
  139. func (db *ObjectDB) BatchGet(ctx SQLContext, objectIDs []cdssdk.ObjectID) ([]cdssdk.Object, error) {
  140. if len(objectIDs) == 0 {
  141. return nil, nil
  142. }
  143. var objs []cdssdk.Object
  144. err := ctx.Table("Object").Where("ObjectID IN ?", objectIDs).Order("ObjectID ASC").Find(&objs).Error
  145. if err != nil {
  146. return nil, err
  147. }
  148. return objs, nil
  149. }
  150. func (db *ObjectDB) BatchGetByPackagePath(ctx SQLContext, pkgID cdssdk.PackageID, pathes []string) ([]cdssdk.Object, error) {
  151. if len(pathes) == 0 {
  152. return nil, nil
  153. }
  154. var objs []cdssdk.Object
  155. err := ctx.Table("Object").Where("PackageID = ? AND Path IN ?", pkgID, pathes).Find(&objs).Error
  156. if err != nil {
  157. return nil, err
  158. }
  159. return objs, nil
  160. }
  161. func (db *ObjectDB) GetDetail(ctx SQLContext, objectID cdssdk.ObjectID) (stgmod.ObjectDetail, error) {
  162. var obj cdssdk.Object
  163. err := ctx.Table("Object").Where("ObjectID = ?", objectID).First(&obj).Error
  164. if err != nil {
  165. return stgmod.ObjectDetail{}, fmt.Errorf("getting object: %w", err)
  166. }
  167. // 获取所有的 ObjectBlock
  168. var allBlocks []stgmod.ObjectBlock
  169. err = ctx.Table("ObjectBlock").Where("ObjectID = ?", objectID).Order("`Index` ASC").Find(&allBlocks).Error
  170. if err != nil {
  171. return stgmod.ObjectDetail{}, fmt.Errorf("getting all object blocks: %w", err)
  172. }
  173. // 获取所有的 PinnedObject
  174. var allPinnedObjs []cdssdk.PinnedObject
  175. err = ctx.Table("PinnedObject").Where("ObjectID = ?", objectID).Order("ObjectID ASC").Find(&allPinnedObjs).Error
  176. if err != nil {
  177. return stgmod.ObjectDetail{}, fmt.Errorf("getting all pinned objects: %w", err)
  178. }
  179. pinnedAt := make([]cdssdk.StorageID, len(allPinnedObjs))
  180. for i, po := range allPinnedObjs {
  181. pinnedAt[i] = po.StorageID
  182. }
  183. return stgmod.ObjectDetail{
  184. Object: obj,
  185. Blocks: allBlocks,
  186. PinnedAt: pinnedAt,
  187. }, nil
  188. }
  189. // 仅返回查询到的对象
  190. func (db *ObjectDB) BatchGetDetails(ctx SQLContext, objectIDs []cdssdk.ObjectID) ([]stgmod.ObjectDetail, error) {
  191. var objs []cdssdk.Object
  192. err := ctx.Table("Object").Where("ObjectID IN ?", objectIDs).Order("ObjectID ASC").Find(&objs).Error
  193. if err != nil {
  194. return nil, err
  195. }
  196. // 获取所有的 ObjectBlock
  197. var allBlocks []stgmod.ObjectBlock
  198. err = ctx.Table("ObjectBlock").Where("ObjectID IN ?", objectIDs).Order("ObjectID, `Index` ASC").Find(&allBlocks).Error
  199. if err != nil {
  200. return nil, err
  201. }
  202. // 获取所有的 PinnedObject
  203. var allPinnedObjs []cdssdk.PinnedObject
  204. err = ctx.Table("PinnedObject").Where("ObjectID IN ?", objectIDs).Order("ObjectID ASC").Find(&allPinnedObjs).Error
  205. if err != nil {
  206. return nil, err
  207. }
  208. details := make([]stgmod.ObjectDetail, len(objs))
  209. for i, obj := range objs {
  210. details[i] = stgmod.ObjectDetail{
  211. Object: obj,
  212. }
  213. }
  214. stgmod.DetailsFillObjectBlocks(details, allBlocks)
  215. stgmod.DetailsFillPinnedAt(details, allPinnedObjs)
  216. return details, nil
  217. }
  218. func (db *ObjectDB) Create(ctx SQLContext, obj cdssdk.Object) (cdssdk.ObjectID, error) {
  219. err := ctx.Table("Object").Create(&obj).Error
  220. if err != nil {
  221. return 0, fmt.Errorf("insert object failed, err: %w", err)
  222. }
  223. return obj.ObjectID, nil
  224. }
  225. // 批量创建对象,创建完成后会填充ObjectID。
  226. func (db *ObjectDB) BatchCreate(ctx SQLContext, objs *[]cdssdk.Object) error {
  227. if len(*objs) == 0 {
  228. return nil
  229. }
  230. return ctx.Table("Object").Create(objs).Error
  231. }
  232. // 批量更新对象所有属性,objs中的对象必须包含ObjectID
  233. func (db *ObjectDB) BatchUpdate(ctx SQLContext, objs []cdssdk.Object) error {
  234. if len(objs) == 0 {
  235. return nil
  236. }
  237. return ctx.Clauses(clause.OnConflict{
  238. Columns: []clause.Column{{Name: "ObjectID"}},
  239. UpdateAll: true,
  240. }).Create(objs).Error
  241. }
  242. // 批量更新对象指定属性,objs中的对象只需设置需要更新的属性即可,但:
  243. // 1. 必须包含ObjectID
  244. // 2. 日期类型属性不能设置为0值
  245. func (db *ObjectDB) BatchUpdateColumns(ctx SQLContext, objs []cdssdk.Object, columns []string) error {
  246. if len(objs) == 0 {
  247. return nil
  248. }
  249. return ctx.Clauses(clause.OnConflict{
  250. Columns: []clause.Column{{Name: "ObjectID"}},
  251. DoUpdates: clause.AssignmentColumns(columns),
  252. }).Create(objs).Error
  253. }
  254. func (db *ObjectDB) GetPackageObjects(ctx SQLContext, packageID cdssdk.PackageID) ([]cdssdk.Object, error) {
  255. var ret []cdssdk.Object
  256. err := ctx.Table("Object").Where("PackageID = ?", packageID).Order("ObjectID ASC").Find(&ret).Error
  257. return ret, err
  258. }
  259. func (db *ObjectDB) GetPackageObjectDetails(ctx SQLContext, packageID cdssdk.PackageID) ([]stgmod.ObjectDetail, error) {
  260. var objs []cdssdk.Object
  261. err := ctx.Table("Object").Where("PackageID = ?", packageID).Order("ObjectID ASC").Find(&objs).Error
  262. if err != nil {
  263. return nil, fmt.Errorf("getting objects: %w", err)
  264. }
  265. // 获取所有的 ObjectBlock
  266. var allBlocks []stgmod.ObjectBlock
  267. err = ctx.Table("ObjectBlock").
  268. Select("ObjectBlock.*").
  269. Joins("JOIN Object ON ObjectBlock.ObjectID = Object.ObjectID").
  270. Where("Object.PackageID = ?", packageID).
  271. Order("ObjectBlock.ObjectID, `Index` ASC").
  272. Find(&allBlocks).Error
  273. if err != nil {
  274. return nil, fmt.Errorf("getting all object blocks: %w", err)
  275. }
  276. // 获取所有的 PinnedObject
  277. var allPinnedObjs []cdssdk.PinnedObject
  278. err = ctx.Table("PinnedObject").
  279. Select("PinnedObject.*").
  280. Joins("JOIN Object ON PinnedObject.ObjectID = Object.ObjectID").
  281. Where("Object.PackageID = ?", packageID).
  282. Order("PinnedObject.ObjectID").
  283. Find(&allPinnedObjs).Error
  284. if err != nil {
  285. return nil, fmt.Errorf("getting all pinned objects: %w", err)
  286. }
  287. details := make([]stgmod.ObjectDetail, len(objs))
  288. for i, obj := range objs {
  289. details[i] = stgmod.ObjectDetail{
  290. Object: obj,
  291. }
  292. }
  293. stgmod.DetailsFillObjectBlocks(details, allBlocks)
  294. stgmod.DetailsFillPinnedAt(details, allPinnedObjs)
  295. return details, nil
  296. }
  297. func (db *ObjectDB) GetObjectsIfAnyBlockOnStorage(ctx SQLContext, stgID cdssdk.StorageID) ([]cdssdk.Object, error) {
  298. var objs []cdssdk.Object
  299. err := ctx.Table("Object").Where("ObjectID IN (SELECT ObjectID FROM ObjectBlock WHERE StorageID = ?)", stgID).Order("ObjectID ASC").Find(&objs).Error
  300. if err != nil {
  301. return nil, fmt.Errorf("getting objects: %w", err)
  302. }
  303. return objs, nil
  304. }
  305. func (db *ObjectDB) BatchAdd(ctx SQLContext, packageID cdssdk.PackageID, adds []coormq.AddObjectEntry) ([]cdssdk.Object, error) {
  306. if len(adds) == 0 {
  307. return nil, nil
  308. }
  309. // 收集所有路径
  310. pathes := make([]string, 0, len(adds))
  311. for _, add := range adds {
  312. pathes = append(pathes, add.Path)
  313. }
  314. // 先查询要更新的对象,不存在也没关系
  315. existsObjs, err := db.BatchGetByPackagePath(ctx, packageID, pathes)
  316. if err != nil {
  317. return nil, fmt.Errorf("batch get object by path: %w", err)
  318. }
  319. existsObjsMap := make(map[string]cdssdk.Object)
  320. for _, obj := range existsObjs {
  321. existsObjsMap[obj.Path] = obj
  322. }
  323. var updatingObjs []cdssdk.Object
  324. var addingObjs []cdssdk.Object
  325. for i := range adds {
  326. o := cdssdk.Object{
  327. PackageID: packageID,
  328. Path: adds[i].Path,
  329. Size: adds[i].Size,
  330. FileHash: adds[i].FileHash,
  331. Redundancy: cdssdk.NewNoneRedundancy(), // 首次上传默认使用不分块的none模式
  332. CreateTime: adds[i].UploadTime,
  333. UpdateTime: adds[i].UploadTime,
  334. }
  335. e, ok := existsObjsMap[adds[i].Path]
  336. if ok {
  337. o.ObjectID = e.ObjectID
  338. o.CreateTime = e.CreateTime
  339. updatingObjs = append(updatingObjs, o)
  340. } else {
  341. addingObjs = append(addingObjs, o)
  342. }
  343. }
  344. // 先进行更新
  345. err = db.BatchUpdate(ctx, updatingObjs)
  346. if err != nil {
  347. return nil, fmt.Errorf("batch update objects: %w", err)
  348. }
  349. // 再执行插入,Create函数插入后会填充ObjectID
  350. err = db.BatchCreate(ctx, &addingObjs)
  351. if err != nil {
  352. return nil, fmt.Errorf("batch create objects: %w", err)
  353. }
  354. // 按照add参数的顺序返回结果
  355. affectedObjsMp := make(map[string]cdssdk.Object)
  356. for _, o := range updatingObjs {
  357. affectedObjsMp[o.Path] = o
  358. }
  359. for _, o := range addingObjs {
  360. affectedObjsMp[o.Path] = o
  361. }
  362. affectedObjs := make([]cdssdk.Object, 0, len(affectedObjsMp))
  363. affectedObjIDs := make([]cdssdk.ObjectID, 0, len(affectedObjsMp))
  364. for i := range adds {
  365. obj := affectedObjsMp[adds[i].Path]
  366. affectedObjs = append(affectedObjs, obj)
  367. affectedObjIDs = append(affectedObjIDs, obj.ObjectID)
  368. }
  369. if len(affectedObjIDs) > 0 {
  370. // 批量删除 ObjectBlock
  371. if err := db.ObjectBlock().BatchDeleteByObjectID(ctx, affectedObjIDs); err != nil {
  372. return nil, fmt.Errorf("batch delete object blocks: %w", err)
  373. }
  374. // 批量删除 PinnedObject
  375. if err := db.PinnedObject().BatchDeleteByObjectID(ctx, affectedObjIDs); err != nil {
  376. return nil, fmt.Errorf("batch delete pinned objects: %w", err)
  377. }
  378. }
  379. // 创建 ObjectBlock
  380. objBlocks := make([]stgmod.ObjectBlock, 0, len(adds))
  381. for i, add := range adds {
  382. for _, stgID := range add.StorageIDs {
  383. objBlocks = append(objBlocks, stgmod.ObjectBlock{
  384. ObjectID: affectedObjIDs[i],
  385. Index: 0,
  386. StorageID: stgID,
  387. FileHash: add.FileHash,
  388. Size: add.Size,
  389. })
  390. }
  391. }
  392. if err := db.ObjectBlock().BatchCreate(ctx, objBlocks); err != nil {
  393. return nil, fmt.Errorf("batch create object blocks: %w", err)
  394. }
  395. // 创建 Cache
  396. caches := make([]model.Cache, 0, len(adds))
  397. for _, add := range adds {
  398. for _, stgID := range add.StorageIDs {
  399. caches = append(caches, model.Cache{
  400. FileHash: add.FileHash,
  401. StorageID: stgID,
  402. CreateTime: time.Now(),
  403. Priority: 0,
  404. })
  405. }
  406. }
  407. if err := db.Cache().BatchCreate(ctx, caches); err != nil {
  408. return nil, fmt.Errorf("batch create caches: %w", err)
  409. }
  410. return affectedObjs, nil
  411. }
  412. func (db *ObjectDB) BatchDelete(ctx SQLContext, ids []cdssdk.ObjectID) error {
  413. if len(ids) == 0 {
  414. return nil
  415. }
  416. return ctx.Table("Object").Where("ObjectID IN ?", ids).Delete(&cdssdk.Object{}).Error
  417. }
  418. func (db *ObjectDB) DeleteInPackage(ctx SQLContext, packageID cdssdk.PackageID) error {
  419. return ctx.Table("Object").Where("PackageID = ?", packageID).Delete(&cdssdk.Object{}).Error
  420. }
  421. func (db *ObjectDB) DeleteByPath(ctx SQLContext, packageID cdssdk.PackageID, path string) error {
  422. return ctx.Table("Object").Where("PackageID = ? AND Path = ?", packageID, path).Delete(&cdssdk.Object{}).Error
  423. }
  424. func (db *ObjectDB) MoveByPrefix(ctx SQLContext, oldPkgID cdssdk.PackageID, oldPrefix string, newPkgID cdssdk.PackageID, newPrefix string) error {
  425. return ctx.Table("Object").Where("PackageID = ? AND Path LIKE ?", oldPkgID, escapeLike("", "%", oldPrefix)).
  426. Updates(map[string]any{
  427. "PackageID": newPkgID,
  428. "Path": gorm.Expr("concat(?, substring(Path, ?))", newPrefix, len(oldPrefix)+1),
  429. }).Error
  430. }

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