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

11 months ago
11 months ago
2 years ago
11 months ago
11 months ago
11 months ago
2 years ago
11 months ago
11 months ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883
  1. package mq
  2. import (
  3. "errors"
  4. "fmt"
  5. "time"
  6. "gitlink.org.cn/cloudream/storage/common/pkgs/db2"
  7. "gitlink.org.cn/cloudream/storage/common/pkgs/db2/model"
  8. "gorm.io/gorm"
  9. "github.com/samber/lo"
  10. "gitlink.org.cn/cloudream/common/consts/errorcode"
  11. "gitlink.org.cn/cloudream/common/pkgs/logger"
  12. "gitlink.org.cn/cloudream/common/pkgs/mq"
  13. cdssdk "gitlink.org.cn/cloudream/common/sdks/storage"
  14. "gitlink.org.cn/cloudream/common/sdks/storage/cdsapi"
  15. "gitlink.org.cn/cloudream/common/utils/sort2"
  16. stgmod "gitlink.org.cn/cloudream/storage/common/models"
  17. coormq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/coordinator"
  18. )
  19. func (svc *Service) GetObjects(msg *coormq.GetObjects) (*coormq.GetObjectsResp, *mq.CodeMessage) {
  20. var ret []*cdssdk.Object
  21. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  22. // TODO 应该检查用户是否有每一个Object所在Package的权限
  23. objs, err := svc.db2.Object().BatchGet(tx, msg.ObjectIDs)
  24. if err != nil {
  25. return err
  26. }
  27. objMp := make(map[cdssdk.ObjectID]cdssdk.Object)
  28. for _, obj := range objs {
  29. objMp[obj.ObjectID] = obj
  30. }
  31. for _, objID := range msg.ObjectIDs {
  32. o, ok := objMp[objID]
  33. if ok {
  34. ret = append(ret, &o)
  35. } else {
  36. ret = append(ret, nil)
  37. }
  38. }
  39. return err
  40. })
  41. if err != nil {
  42. logger.WithField("UserID", msg.UserID).
  43. Warn(err.Error())
  44. return nil, mq.Failed(errorcode.OperationFailed, "get objects failed")
  45. }
  46. return mq.ReplyOK(coormq.RespGetObjects(ret))
  47. }
  48. func (svc *Service) GetObjectsByPath(msg *coormq.GetObjectsByPath) (*coormq.GetObjectsByPathResp, *mq.CodeMessage) {
  49. var coms []string
  50. var objs []cdssdk.Object
  51. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  52. var err error
  53. _, err = svc.db2.Package().GetUserPackage(tx, msg.UserID, msg.PackageID)
  54. if err != nil {
  55. return fmt.Errorf("getting package by id: %w", err)
  56. }
  57. if !msg.IsPrefix {
  58. obj, err := svc.db2.Object().GetByPath(tx, msg.PackageID, msg.Path)
  59. if err != nil {
  60. return fmt.Errorf("getting object by path: %w", err)
  61. }
  62. objs = append(objs, obj)
  63. return nil
  64. }
  65. if !msg.NoRecursive {
  66. objs, err = svc.db2.Object().GetWithPathPrefix(tx, msg.PackageID, msg.Path)
  67. if err != nil {
  68. return fmt.Errorf("getting objects with prefix: %w", err)
  69. }
  70. return nil
  71. }
  72. coms, err = svc.db2.Object().GetCommonPrefixes(tx, msg.PackageID, msg.Path)
  73. if err != nil {
  74. return fmt.Errorf("getting common prefixes: %w", err)
  75. }
  76. objs, err = svc.db2.Object().GetDirectChildren(tx, msg.PackageID, msg.Path)
  77. if err != nil {
  78. return fmt.Errorf("getting direct children: %w", err)
  79. }
  80. return nil
  81. })
  82. if err != nil {
  83. logger.WithField("PathPrefix", msg.Path).Warn(err.Error())
  84. return nil, mq.Failed(errorcode.OperationFailed, "get objects with prefix failed")
  85. }
  86. return mq.ReplyOK(coormq.RespGetObjectsByPath(coms, objs))
  87. }
  88. func (svc *Service) GetPackageObjects(msg *coormq.GetPackageObjects) (*coormq.GetPackageObjectsResp, *mq.CodeMessage) {
  89. var objs []cdssdk.Object
  90. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  91. _, err := svc.db2.Package().GetUserPackage(tx, msg.UserID, msg.PackageID)
  92. if err != nil {
  93. return fmt.Errorf("getting package by id: %w", err)
  94. }
  95. objs, err = svc.db2.Object().GetPackageObjects(tx, msg.PackageID)
  96. if err != nil {
  97. return fmt.Errorf("getting package objects: %w", err)
  98. }
  99. return nil
  100. })
  101. if err != nil {
  102. logger.WithField("UserID", msg.UserID).WithField("PackageID", msg.PackageID).
  103. Warn(err.Error())
  104. return nil, mq.Failed(errorcode.OperationFailed, "get package objects failed")
  105. }
  106. return mq.ReplyOK(coormq.RespGetPackageObjects(objs))
  107. }
  108. func (svc *Service) GetPackageObjectDetails(msg *coormq.GetPackageObjectDetails) (*coormq.GetPackageObjectDetailsResp, *mq.CodeMessage) {
  109. var details []stgmod.ObjectDetail
  110. // 必须放在事务里进行,因为GetPackageBlockDetails是由多次数据库操作组成,必须保证数据的一致性
  111. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  112. var err error
  113. _, err = svc.db2.Package().GetByID(tx, msg.PackageID)
  114. if err != nil {
  115. return fmt.Errorf("getting package by id: %w", err)
  116. }
  117. details, err = svc.db2.Object().GetPackageObjectDetails(tx, msg.PackageID)
  118. if err != nil {
  119. return fmt.Errorf("getting package block details: %w", err)
  120. }
  121. return nil
  122. })
  123. if err != nil {
  124. logger.WithField("PackageID", msg.PackageID).Warn(err.Error())
  125. return nil, mq.Failed(errorcode.OperationFailed, "get package object block details failed")
  126. }
  127. return mq.ReplyOK(coormq.RespPackageObjectDetails(details))
  128. }
  129. func (svc *Service) GetObjectDetails(msg *coormq.GetObjectDetails) (*coormq.GetObjectDetailsResp, *mq.CodeMessage) {
  130. detailsMp := make(map[cdssdk.ObjectID]*stgmod.ObjectDetail)
  131. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  132. var err error
  133. msg.ObjectIDs = sort2.SortAsc(msg.ObjectIDs)
  134. // 根据ID依次查询Object,ObjectBlock,PinnedObject,并根据升序的特点进行合并
  135. objs, err := svc.db2.Object().BatchGet(tx, msg.ObjectIDs)
  136. if err != nil {
  137. return fmt.Errorf("batch get objects: %w", err)
  138. }
  139. for _, obj := range objs {
  140. detailsMp[obj.ObjectID] = &stgmod.ObjectDetail{
  141. Object: obj,
  142. }
  143. }
  144. // 查询合并
  145. blocks, err := svc.db2.ObjectBlock().BatchGetByObjectID(tx, msg.ObjectIDs)
  146. if err != nil {
  147. return fmt.Errorf("batch get object blocks: %w", err)
  148. }
  149. for _, block := range blocks {
  150. d := detailsMp[block.ObjectID]
  151. d.Blocks = append(d.Blocks, block)
  152. }
  153. // 查询合并
  154. pinneds, err := svc.db2.PinnedObject().BatchGetByObjectID(tx, msg.ObjectIDs)
  155. if err != nil {
  156. return fmt.Errorf("batch get pinned objects: %w", err)
  157. }
  158. for _, pinned := range pinneds {
  159. d := detailsMp[pinned.ObjectID]
  160. d.PinnedAt = append(d.PinnedAt, pinned.StorageID)
  161. }
  162. return nil
  163. })
  164. if err != nil {
  165. logger.Warn(err.Error())
  166. return nil, mq.Failed(errorcode.OperationFailed, "get object details failed")
  167. }
  168. details := make([]*stgmod.ObjectDetail, len(msg.ObjectIDs))
  169. for i, objID := range msg.ObjectIDs {
  170. details[i] = detailsMp[objID]
  171. }
  172. return mq.ReplyOK(coormq.RespGetObjectDetails(details))
  173. }
  174. func (svc *Service) UpdateObjectRedundancy(msg *coormq.UpdateObjectRedundancy) (*coormq.UpdateObjectRedundancyResp, *mq.CodeMessage) {
  175. err := svc.db2.DoTx(func(ctx db2.SQLContext) error {
  176. db := svc.db2
  177. objs := msg.Updatings
  178. nowTime := time.Now()
  179. objIDs := make([]cdssdk.ObjectID, 0, len(objs))
  180. for _, obj := range objs {
  181. objIDs = append(objIDs, obj.ObjectID)
  182. }
  183. avaiIDs, err := db.Object().BatchTestObjectID(ctx, objIDs)
  184. if err != nil {
  185. return fmt.Errorf("batch test object id: %w", err)
  186. }
  187. // 过滤掉已经不存在的对象。
  188. // 注意,objIDs没有被过滤,因为后续逻辑不过滤也不会出错
  189. objs = lo.Filter(objs, func(obj coormq.UpdatingObjectRedundancy, _ int) bool {
  190. return avaiIDs[obj.ObjectID]
  191. })
  192. dummyObjs := make([]cdssdk.Object, 0, len(objs))
  193. for _, obj := range objs {
  194. dummyObjs = append(dummyObjs, cdssdk.Object{
  195. ObjectID: obj.ObjectID,
  196. FileHash: obj.FileHash,
  197. Size: obj.Size,
  198. Redundancy: obj.Redundancy,
  199. CreateTime: nowTime, // 实际不会更新,只因为不能是0值
  200. UpdateTime: nowTime,
  201. })
  202. }
  203. err = db.Object().BatchUpdateColumns(ctx, dummyObjs, []string{"FileHash", "Size", "Redundancy", "UpdateTime"})
  204. if err != nil {
  205. return fmt.Errorf("batch update object redundancy: %w", err)
  206. }
  207. // 删除原本所有的编码块记录,重新添加
  208. err = db.ObjectBlock().BatchDeleteByObjectID(ctx, objIDs)
  209. if err != nil {
  210. return fmt.Errorf("batch delete object blocks: %w", err)
  211. }
  212. // 删除原本Pin住的Object。暂不考虑FileHash没有变化的情况
  213. err = db.PinnedObject().BatchDeleteByObjectID(ctx, objIDs)
  214. if err != nil {
  215. return fmt.Errorf("batch delete pinned object: %w", err)
  216. }
  217. blocks := make([]stgmod.ObjectBlock, 0, len(objs))
  218. for _, obj := range objs {
  219. blocks = append(blocks, obj.Blocks...)
  220. }
  221. err = db.ObjectBlock().BatchCreate(ctx, blocks)
  222. if err != nil {
  223. return fmt.Errorf("batch create object blocks: %w", err)
  224. }
  225. caches := make([]model.Cache, 0, len(objs))
  226. for _, obj := range objs {
  227. for _, blk := range obj.Blocks {
  228. caches = append(caches, model.Cache{
  229. FileHash: blk.FileHash,
  230. StorageID: blk.StorageID,
  231. CreateTime: nowTime,
  232. Priority: 0,
  233. })
  234. }
  235. }
  236. err = db.Cache().BatchCreate(ctx, caches)
  237. if err != nil {
  238. return fmt.Errorf("batch create object caches: %w", err)
  239. }
  240. pinneds := make([]cdssdk.PinnedObject, 0, len(objs))
  241. for _, obj := range objs {
  242. for _, p := range obj.PinnedAt {
  243. pinneds = append(pinneds, cdssdk.PinnedObject{
  244. ObjectID: obj.ObjectID,
  245. StorageID: p,
  246. CreateTime: nowTime,
  247. })
  248. }
  249. }
  250. err = db.PinnedObject().BatchTryCreate(ctx, pinneds)
  251. if err != nil {
  252. return fmt.Errorf("batch create pinned objects: %w", err)
  253. }
  254. return nil
  255. })
  256. if err != nil {
  257. logger.Warnf("batch updating redundancy: %s", err.Error())
  258. return nil, mq.Failed(errorcode.OperationFailed, "batch update redundancy failed")
  259. }
  260. return mq.ReplyOK(coormq.RespUpdateObjectRedundancy())
  261. }
  262. func (svc *Service) UpdateObjectInfos(msg *coormq.UpdateObjectInfos) (*coormq.UpdateObjectInfosResp, *mq.CodeMessage) {
  263. var sucs []cdssdk.ObjectID
  264. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  265. msg.Updatings = sort2.Sort(msg.Updatings, func(o1, o2 cdsapi.UpdatingObject) int {
  266. return sort2.Cmp(o1.ObjectID, o2.ObjectID)
  267. })
  268. objIDs := make([]cdssdk.ObjectID, len(msg.Updatings))
  269. for i, obj := range msg.Updatings {
  270. objIDs[i] = obj.ObjectID
  271. }
  272. oldObjs, err := svc.db2.Object().BatchGet(tx, objIDs)
  273. if err != nil {
  274. return fmt.Errorf("batch getting objects: %w", err)
  275. }
  276. oldObjIDs := make([]cdssdk.ObjectID, len(oldObjs))
  277. for i, obj := range oldObjs {
  278. oldObjIDs[i] = obj.ObjectID
  279. }
  280. avaiUpdatings, notExistsObjs := pickByObjectIDs(msg.Updatings, oldObjIDs, func(obj cdsapi.UpdatingObject) cdssdk.ObjectID { return obj.ObjectID })
  281. if len(notExistsObjs) > 0 {
  282. // TODO 部分对象已经不存在
  283. }
  284. newObjs := make([]cdssdk.Object, len(avaiUpdatings))
  285. for i := range newObjs {
  286. newObjs[i] = oldObjs[i]
  287. avaiUpdatings[i].ApplyTo(&newObjs[i])
  288. }
  289. err = svc.db2.Object().BatchUpdate(tx, newObjs)
  290. if err != nil {
  291. return fmt.Errorf("batch create or update: %w", err)
  292. }
  293. sucs = lo.Map(newObjs, func(obj cdssdk.Object, _ int) cdssdk.ObjectID { return obj.ObjectID })
  294. return nil
  295. })
  296. if err != nil {
  297. logger.Warnf("batch updating objects: %s", err.Error())
  298. return nil, mq.Failed(errorcode.OperationFailed, "batch update objects failed")
  299. }
  300. return mq.ReplyOK(coormq.RespUpdateObjectInfos(sucs))
  301. }
  302. // 根据objIDs从objs中挑选Object。
  303. // len(objs) >= len(objIDs)
  304. func pickByObjectIDs[T any](objs []T, objIDs []cdssdk.ObjectID, getID func(T) cdssdk.ObjectID) (picked []T, notFound []T) {
  305. objIdx := 0
  306. idIdx := 0
  307. for idIdx < len(objIDs) && objIdx < len(objs) {
  308. if getID(objs[objIdx]) < objIDs[idIdx] {
  309. notFound = append(notFound, objs[objIdx])
  310. objIdx++
  311. continue
  312. }
  313. picked = append(picked, objs[objIdx])
  314. objIdx++
  315. idIdx++
  316. }
  317. return
  318. }
  319. func (svc *Service) MoveObjects(msg *coormq.MoveObjects) (*coormq.MoveObjectsResp, *mq.CodeMessage) {
  320. var sucs []cdssdk.ObjectID
  321. var evt []*stgmod.BodyObjectInfoUpdated
  322. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  323. msg.Movings = sort2.Sort(msg.Movings, func(o1, o2 cdsapi.MovingObject) int {
  324. return sort2.Cmp(o1.ObjectID, o2.ObjectID)
  325. })
  326. objIDs := make([]cdssdk.ObjectID, len(msg.Movings))
  327. for i, obj := range msg.Movings {
  328. objIDs[i] = obj.ObjectID
  329. }
  330. oldObjs, err := svc.db2.Object().BatchGet(tx, objIDs)
  331. if err != nil {
  332. return fmt.Errorf("batch getting objects: %w", err)
  333. }
  334. oldObjIDs := make([]cdssdk.ObjectID, len(oldObjs))
  335. for i, obj := range oldObjs {
  336. oldObjIDs[i] = obj.ObjectID
  337. }
  338. // 找出仍在数据库的Object
  339. avaiMovings, notExistsObjs := pickByObjectIDs(msg.Movings, oldObjIDs, func(obj cdsapi.MovingObject) cdssdk.ObjectID { return obj.ObjectID })
  340. if len(notExistsObjs) > 0 {
  341. // TODO 部分对象已经不存在
  342. }
  343. // 筛选出PackageID变化、Path变化的对象,这两种对象要检测改变后是否有冲突
  344. var pkgIDChangedObjs []cdssdk.Object
  345. var pathChangedObjs []cdssdk.Object
  346. for i := range avaiMovings {
  347. if avaiMovings[i].PackageID != oldObjs[i].PackageID {
  348. newObj := oldObjs[i]
  349. avaiMovings[i].ApplyTo(&newObj)
  350. pkgIDChangedObjs = append(pkgIDChangedObjs, newObj)
  351. } else if avaiMovings[i].Path != oldObjs[i].Path {
  352. newObj := oldObjs[i]
  353. avaiMovings[i].ApplyTo(&newObj)
  354. pathChangedObjs = append(pathChangedObjs, newObj)
  355. }
  356. }
  357. var newObjs []cdssdk.Object
  358. // 对于PackageID发生变化的对象,需要检查目标Package内是否存在同Path的对象
  359. checkedObjs, err := svc.checkPackageChangedObjects(tx, msg.UserID, pkgIDChangedObjs)
  360. if err != nil {
  361. return err
  362. }
  363. newObjs = append(newObjs, checkedObjs...)
  364. // 对于只有Path发生变化的对象,则检查同Package内有没有同Path的对象
  365. checkedObjs, err = svc.checkPathChangedObjects(tx, msg.UserID, pathChangedObjs)
  366. if err != nil {
  367. return err
  368. }
  369. newObjs = append(newObjs, checkedObjs...)
  370. err = svc.db2.Object().BatchUpdate(tx, newObjs)
  371. if err != nil {
  372. return fmt.Errorf("batch create or update: %w", err)
  373. }
  374. sucs = lo.Map(newObjs, func(obj cdssdk.Object, _ int) cdssdk.ObjectID { return obj.ObjectID })
  375. evt = lo.Map(newObjs, func(obj cdssdk.Object, _ int) *stgmod.BodyObjectInfoUpdated {
  376. return &stgmod.BodyObjectInfoUpdated{
  377. Object: obj,
  378. }
  379. })
  380. return nil
  381. })
  382. if err != nil {
  383. logger.Warn(err.Error())
  384. return nil, mq.Failed(errorcode.OperationFailed, "move objects failed")
  385. }
  386. for _, e := range evt {
  387. svc.evtPub.Publish(e)
  388. }
  389. return mq.ReplyOK(coormq.RespMoveObjects(sucs))
  390. }
  391. func (svc *Service) checkPackageChangedObjects(tx db2.SQLContext, userID cdssdk.UserID, objs []cdssdk.Object) ([]cdssdk.Object, error) {
  392. if len(objs) == 0 {
  393. return nil, nil
  394. }
  395. type PackageObjects struct {
  396. PackageID cdssdk.PackageID
  397. ObjectByPath map[string]*cdssdk.Object
  398. }
  399. packages := make(map[cdssdk.PackageID]*PackageObjects)
  400. for _, obj := range objs {
  401. pkg, ok := packages[obj.PackageID]
  402. if !ok {
  403. pkg = &PackageObjects{
  404. PackageID: obj.PackageID,
  405. ObjectByPath: make(map[string]*cdssdk.Object),
  406. }
  407. packages[obj.PackageID] = pkg
  408. }
  409. if pkg.ObjectByPath[obj.Path] == nil {
  410. o := obj
  411. pkg.ObjectByPath[obj.Path] = &o
  412. } else {
  413. // TODO 有两个对象移动到同一个路径,有冲突
  414. }
  415. }
  416. var willUpdateObjs []cdssdk.Object
  417. for _, pkg := range packages {
  418. _, err := svc.db2.Package().GetUserPackage(tx, userID, pkg.PackageID)
  419. if errors.Is(err, gorm.ErrRecordNotFound) {
  420. continue
  421. }
  422. if err != nil {
  423. return nil, fmt.Errorf("getting user package by id: %w", err)
  424. }
  425. existsObjs, err := svc.db2.Object().BatchGetByPackagePath(tx, pkg.PackageID, lo.Keys(pkg.ObjectByPath))
  426. if err != nil {
  427. return nil, fmt.Errorf("batch getting objects by package path: %w", err)
  428. }
  429. // 标记冲突的对象
  430. for _, obj := range existsObjs {
  431. pkg.ObjectByPath[obj.Path] = nil
  432. // TODO 目标Package内有冲突的对象
  433. }
  434. for _, obj := range pkg.ObjectByPath {
  435. if obj == nil {
  436. continue
  437. }
  438. willUpdateObjs = append(willUpdateObjs, *obj)
  439. }
  440. }
  441. return willUpdateObjs, nil
  442. }
  443. func (svc *Service) checkPathChangedObjects(tx db2.SQLContext, userID cdssdk.UserID, objs []cdssdk.Object) ([]cdssdk.Object, error) {
  444. if len(objs) == 0 {
  445. return nil, nil
  446. }
  447. objByPath := make(map[string]*cdssdk.Object)
  448. for _, obj := range objs {
  449. if objByPath[obj.Path] == nil {
  450. o := obj
  451. objByPath[obj.Path] = &o
  452. } else {
  453. // TODO 有两个对象移动到同一个路径,有冲突
  454. }
  455. }
  456. _, err := svc.db2.Package().GetUserPackage(tx, userID, objs[0].PackageID)
  457. if errors.Is(err, gorm.ErrRecordNotFound) {
  458. return nil, nil
  459. }
  460. if err != nil {
  461. return nil, fmt.Errorf("getting user package by id: %w", err)
  462. }
  463. existsObjs, err := svc.db2.Object().BatchGetByPackagePath(tx, objs[0].PackageID, lo.Map(objs, func(obj cdssdk.Object, idx int) string { return obj.Path }))
  464. if err != nil {
  465. return nil, fmt.Errorf("batch getting objects by package path: %w", err)
  466. }
  467. // 不支持两个对象交换位置的情况,因为数据库不支持
  468. for _, obj := range existsObjs {
  469. objByPath[obj.Path] = nil
  470. }
  471. var willMoveObjs []cdssdk.Object
  472. for _, obj := range objByPath {
  473. if obj == nil {
  474. continue
  475. }
  476. willMoveObjs = append(willMoveObjs, *obj)
  477. }
  478. return willMoveObjs, nil
  479. }
  480. func (svc *Service) DeleteObjects(msg *coormq.DeleteObjects) (*coormq.DeleteObjectsResp, *mq.CodeMessage) {
  481. var sucs []cdssdk.ObjectID
  482. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  483. avaiIDs, err := svc.db2.Object().BatchTestObjectID(tx, msg.ObjectIDs)
  484. if err != nil {
  485. return fmt.Errorf("batch testing object id: %w", err)
  486. }
  487. sucs = lo.Keys(avaiIDs)
  488. err = svc.db2.Object().BatchDelete(tx, msg.ObjectIDs)
  489. if err != nil {
  490. return fmt.Errorf("batch deleting objects: %w", err)
  491. }
  492. err = svc.db2.ObjectBlock().BatchDeleteByObjectID(tx, msg.ObjectIDs)
  493. if err != nil {
  494. return fmt.Errorf("batch deleting object blocks: %w", err)
  495. }
  496. err = svc.db2.PinnedObject().BatchDeleteByObjectID(tx, msg.ObjectIDs)
  497. if err != nil {
  498. return fmt.Errorf("batch deleting pinned objects: %w", err)
  499. }
  500. err = svc.db2.ObjectAccessStat().BatchDeleteByObjectID(tx, msg.ObjectIDs)
  501. if err != nil {
  502. return fmt.Errorf("batch deleting object access stats: %w", err)
  503. }
  504. return nil
  505. })
  506. if err != nil {
  507. logger.Warnf("batch deleting objects: %s", err.Error())
  508. return nil, mq.Failed(errorcode.OperationFailed, "batch delete objects failed")
  509. }
  510. for _, objID := range sucs {
  511. svc.evtPub.Publish(&stgmod.BodyObjectDeleted{
  512. ObjectID: objID,
  513. })
  514. }
  515. return mq.ReplyOK(coormq.RespDeleteObjects(sucs))
  516. }
  517. func (svc *Service) CloneObjects(msg *coormq.CloneObjects) (*coormq.CloneObjectsResp, *mq.CodeMessage) {
  518. type CloningObject struct {
  519. Cloning cdsapi.CloningObject
  520. OrgIndex int
  521. }
  522. type PackageClonings struct {
  523. PackageID cdssdk.PackageID
  524. Clonings map[string]CloningObject
  525. }
  526. var evt []*stgmod.BodyNewOrUpdateObject
  527. // TODO 要检查用户是否有Object、Package的权限
  528. clonings := make(map[cdssdk.PackageID]*PackageClonings)
  529. for i, cloning := range msg.Clonings {
  530. pkg, ok := clonings[cloning.NewPackageID]
  531. if !ok {
  532. pkg = &PackageClonings{
  533. PackageID: cloning.NewPackageID,
  534. Clonings: make(map[string]CloningObject),
  535. }
  536. clonings[cloning.NewPackageID] = pkg
  537. }
  538. pkg.Clonings[cloning.NewPath] = CloningObject{
  539. Cloning: cloning,
  540. OrgIndex: i,
  541. }
  542. }
  543. ret := make([]*cdssdk.Object, len(msg.Clonings))
  544. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  545. // 剔除掉新路径已经存在的对象
  546. for _, pkg := range clonings {
  547. exists, err := svc.db2.Object().BatchGetByPackagePath(tx, pkg.PackageID, lo.Keys(pkg.Clonings))
  548. if err != nil {
  549. return fmt.Errorf("batch getting objects by package path: %w", err)
  550. }
  551. for _, obj := range exists {
  552. delete(pkg.Clonings, obj.Path)
  553. }
  554. }
  555. // 删除目的Package不存在的对象
  556. newPkg, err := svc.db2.Package().BatchTestPackageID(tx, lo.Keys(clonings))
  557. if err != nil {
  558. return fmt.Errorf("batch testing package id: %w", err)
  559. }
  560. for _, pkg := range clonings {
  561. if !newPkg[pkg.PackageID] {
  562. delete(clonings, pkg.PackageID)
  563. }
  564. }
  565. var avaiClonings []CloningObject
  566. var avaiObjIDs []cdssdk.ObjectID
  567. for _, pkg := range clonings {
  568. for _, cloning := range pkg.Clonings {
  569. avaiClonings = append(avaiClonings, cloning)
  570. avaiObjIDs = append(avaiObjIDs, cloning.Cloning.ObjectID)
  571. }
  572. }
  573. avaiDetails, err := svc.db2.Object().BatchGetDetails(tx, avaiObjIDs)
  574. if err != nil {
  575. return fmt.Errorf("batch getting object details: %w", err)
  576. }
  577. avaiDetailsMap := make(map[cdssdk.ObjectID]stgmod.ObjectDetail)
  578. for _, detail := range avaiDetails {
  579. avaiDetailsMap[detail.Object.ObjectID] = detail
  580. }
  581. oldAvaiClonings := avaiClonings
  582. avaiClonings = nil
  583. var newObjs []cdssdk.Object
  584. for _, cloning := range oldAvaiClonings {
  585. // 进一步剔除原始对象不存在的情况
  586. detail, ok := avaiDetailsMap[cloning.Cloning.ObjectID]
  587. if !ok {
  588. continue
  589. }
  590. avaiClonings = append(avaiClonings, cloning)
  591. newObj := detail.Object
  592. newObj.ObjectID = 0
  593. newObj.Path = cloning.Cloning.NewPath
  594. newObj.PackageID = cloning.Cloning.NewPackageID
  595. newObjs = append(newObjs, newObj)
  596. }
  597. // 先创建出新对象
  598. err = svc.db2.Object().BatchCreate(tx, &newObjs)
  599. if err != nil {
  600. return fmt.Errorf("batch creating objects: %w", err)
  601. }
  602. // 创建了新对象就能拿到新对象ID,再创建新对象块
  603. var newBlks []stgmod.ObjectBlock
  604. for i, cloning := range avaiClonings {
  605. oldBlks := avaiDetailsMap[cloning.Cloning.ObjectID].Blocks
  606. for _, blk := range oldBlks {
  607. newBlk := blk
  608. newBlk.ObjectID = newObjs[i].ObjectID
  609. newBlks = append(newBlks, newBlk)
  610. }
  611. }
  612. err = svc.db2.ObjectBlock().BatchCreate(tx, newBlks)
  613. if err != nil {
  614. return fmt.Errorf("batch creating object blocks: %w", err)
  615. }
  616. for i, cloning := range avaiClonings {
  617. ret[cloning.OrgIndex] = &newObjs[i]
  618. }
  619. for i, cloning := range avaiClonings {
  620. var evtBlks []stgmod.BlockDistributionObjectInfo
  621. blkType := getBlockTypeFromRed(newObjs[i].Redundancy)
  622. oldBlks := avaiDetailsMap[cloning.Cloning.ObjectID].Blocks
  623. for _, blk := range oldBlks {
  624. evtBlks = append(evtBlks, stgmod.BlockDistributionObjectInfo{
  625. BlockType: blkType,
  626. Index: blk.Index,
  627. StorageID: blk.StorageID,
  628. })
  629. }
  630. evt = append(evt, &stgmod.BodyNewOrUpdateObject{
  631. Info: newObjs[i],
  632. BlockDistribution: evtBlks,
  633. })
  634. }
  635. return nil
  636. })
  637. if err != nil {
  638. logger.Warnf("cloning objects: %s", err.Error())
  639. return nil, mq.Failed(errorcode.OperationFailed, err.Error())
  640. }
  641. for _, e := range evt {
  642. svc.evtPub.Publish(e)
  643. }
  644. return mq.ReplyOK(coormq.RespCloneObjects(ret))
  645. }
  646. func (svc *Service) NewMultipartUploadObject(msg *coormq.NewMultipartUploadObject) (*coormq.NewMultipartUploadObjectResp, *mq.CodeMessage) {
  647. var obj cdssdk.Object
  648. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  649. oldObj, err := svc.db2.Object().GetByPath(tx, msg.PackageID, msg.Path)
  650. if err == nil {
  651. obj = oldObj
  652. err := svc.db2.ObjectBlock().DeleteByObjectID(tx, obj.ObjectID)
  653. if err != nil {
  654. return fmt.Errorf("delete object blocks: %w", err)
  655. }
  656. obj.FileHash = cdssdk.EmptyHash
  657. obj.Size = 0
  658. obj.Redundancy = cdssdk.NewMultipartUploadRedundancy()
  659. obj.UpdateTime = time.Now()
  660. err = svc.db2.Object().BatchUpdate(tx, []cdssdk.Object{obj})
  661. if err != nil {
  662. return fmt.Errorf("update object: %w", err)
  663. }
  664. return nil
  665. }
  666. obj = cdssdk.Object{
  667. PackageID: msg.PackageID,
  668. Path: msg.Path,
  669. FileHash: cdssdk.EmptyHash,
  670. Size: 0,
  671. Redundancy: cdssdk.NewMultipartUploadRedundancy(),
  672. CreateTime: time.Now(),
  673. UpdateTime: time.Now(),
  674. }
  675. objID, err := svc.db2.Object().Create(tx, obj)
  676. if err != nil {
  677. return fmt.Errorf("create object: %w", err)
  678. }
  679. obj.ObjectID = objID
  680. return nil
  681. })
  682. if err != nil {
  683. logger.Warnf("new multipart upload object: %s", err.Error())
  684. return nil, mq.Failed(errorcode.OperationFailed, fmt.Sprintf("new multipart upload object: %v", err))
  685. }
  686. return mq.ReplyOK(coormq.RespNewMultipartUploadObject(obj))
  687. }
  688. func (svc *Service) AddMultipartUploadPart(msg *coormq.AddMultipartUploadPart) (*coormq.AddMultipartUploadPartResp, *mq.CodeMessage) {
  689. err := svc.db2.DoTx(func(tx db2.SQLContext) error {
  690. obj, err := svc.db2.Object().GetByID(tx, msg.ObjectID)
  691. if err != nil {
  692. return fmt.Errorf("getting object by id: %w", err)
  693. }
  694. _, ok := obj.Redundancy.(*cdssdk.MultipartUploadRedundancy)
  695. if !ok {
  696. return fmt.Errorf("object is not a multipart upload object")
  697. }
  698. blks, err := svc.db2.ObjectBlock().BatchGetByObjectID(tx, []cdssdk.ObjectID{obj.ObjectID})
  699. if err != nil {
  700. return fmt.Errorf("batch getting object blocks: %w", err)
  701. }
  702. blks = lo.Reject(blks, func(blk stgmod.ObjectBlock, idx int) bool { return blk.Index == msg.Block.Index })
  703. blks = append(blks, msg.Block)
  704. blks = sort2.Sort(blks, func(a, b stgmod.ObjectBlock) int { return a.Index - b.Index })
  705. totalSize := int64(0)
  706. var hashes [][]byte
  707. for _, blk := range blks {
  708. totalSize += blk.Size
  709. hashes = append(hashes, blk.FileHash.GetHashBytes())
  710. }
  711. newObjHash := cdssdk.CalculateCompositeHash(hashes)
  712. obj.Size = totalSize
  713. obj.FileHash = newObjHash
  714. obj.UpdateTime = time.Now()
  715. err = svc.db2.ObjectBlock().DeleteByObjectIDIndex(tx, msg.ObjectID, msg.Block.Index)
  716. if err != nil {
  717. return fmt.Errorf("delete object block: %w", err)
  718. }
  719. err = svc.db2.ObjectBlock().Create(tx, msg.ObjectID, msg.Block.Index, msg.Block.StorageID, msg.Block.FileHash, msg.Block.Size)
  720. if err != nil {
  721. return fmt.Errorf("create object block: %w", err)
  722. }
  723. err = svc.db2.Object().BatchUpdate(tx, []cdssdk.Object{obj})
  724. if err != nil {
  725. return fmt.Errorf("update object: %w", err)
  726. }
  727. return nil
  728. })
  729. if err != nil {
  730. logger.Warnf("add multipart upload part: %s", err.Error())
  731. code := errorcode.OperationFailed
  732. if errors.Is(err, gorm.ErrRecordNotFound) {
  733. code = errorcode.DataNotFound
  734. }
  735. return nil, mq.Failed(code, fmt.Sprintf("add multipart upload part: %v", err))
  736. }
  737. return mq.ReplyOK(coormq.RespAddMultipartUploadPart())
  738. }

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