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.

package.go 5.9 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
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229
  1. package cmdline
  2. import (
  3. "fmt"
  4. "io"
  5. "os"
  6. "path/filepath"
  7. "time"
  8. "github.com/jedib0t/go-pretty/v6/table"
  9. cdssdk "gitlink.org.cn/cloudream/common/sdks/storage"
  10. "gitlink.org.cn/cloudream/storage/common/pkgs/iterator"
  11. )
  12. func PackageListBucketPackages(ctx CommandContext, bucketID cdssdk.BucketID) error {
  13. userID := cdssdk.UserID(0)
  14. packages, err := ctx.Cmdline.Svc.BucketSvc().GetBucketPackages(userID, bucketID)
  15. if err != nil {
  16. return err
  17. }
  18. fmt.Printf("Find %d packages in bucket %d for user %d:\n", len(packages), bucketID, userID)
  19. tb := table.NewWriter()
  20. tb.AppendHeader(table.Row{"ID", "Name", "BucketID", "State"})
  21. for _, obj := range packages {
  22. tb.AppendRow(table.Row{obj.PackageID, obj.Name, obj.BucketID, obj.State})
  23. }
  24. fmt.Println(tb.Render())
  25. return nil
  26. }
  27. func PackageDownloadPackage(ctx CommandContext, outputDir string, packageID cdssdk.PackageID) error {
  28. err := os.MkdirAll(outputDir, os.ModePerm)
  29. if err != nil {
  30. return fmt.Errorf("create output directory %s failed, err: %w", outputDir, err)
  31. }
  32. // 下载文件
  33. objIter, err := ctx.Cmdline.Svc.PackageSvc().DownloadPackage(0, packageID)
  34. if err != nil {
  35. return fmt.Errorf("download object failed, err: %w", err)
  36. }
  37. defer objIter.Close()
  38. for {
  39. objInfo, err := objIter.MoveNext()
  40. if err == iterator.ErrNoMoreItem {
  41. break
  42. }
  43. if err != nil {
  44. return err
  45. }
  46. err = func() error {
  47. defer objInfo.File.Close()
  48. fullPath := filepath.Join(outputDir, objInfo.Object.Path)
  49. dirPath := filepath.Dir(fullPath)
  50. if err := os.MkdirAll(dirPath, 0755); err != nil {
  51. return fmt.Errorf("creating object dir: %w", err)
  52. }
  53. outputFile, err := os.Create(fullPath)
  54. if err != nil {
  55. return fmt.Errorf("creating object file: %w", err)
  56. }
  57. defer outputFile.Close()
  58. _, err = io.Copy(outputFile, objInfo.File)
  59. if err != nil {
  60. return fmt.Errorf("copy object data to local file failed, err: %w", err)
  61. }
  62. return nil
  63. }()
  64. if err != nil {
  65. return err
  66. }
  67. }
  68. return nil
  69. }
  70. func PackageCreatePackage(ctx CommandContext, rootPath string, bucketID cdssdk.BucketID, name string, nodeAffinity []cdssdk.NodeID) error {
  71. rootPath = filepath.Clean(rootPath)
  72. var uploadFilePathes []string
  73. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  74. if err != nil {
  75. return nil
  76. }
  77. if !fi.IsDir() {
  78. uploadFilePathes = append(uploadFilePathes, fname)
  79. }
  80. return nil
  81. })
  82. if err != nil {
  83. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  84. }
  85. var nodeAff *cdssdk.NodeID
  86. if len(nodeAffinity) > 0 {
  87. n := cdssdk.NodeID(nodeAffinity[0])
  88. nodeAff = &n
  89. }
  90. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  91. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartCreatingPackage(0, bucketID, name, objIter, nodeAff)
  92. if err != nil {
  93. return fmt.Errorf("upload file data failed, err: %w", err)
  94. }
  95. for {
  96. complete, uploadObjectResult, err := ctx.Cmdline.Svc.PackageSvc().WaitCreatingPackage(taskID, time.Second*5)
  97. if complete {
  98. if err != nil {
  99. return fmt.Errorf("uploading package: %w", err)
  100. }
  101. tb := table.NewWriter()
  102. tb.AppendHeader(table.Row{"Path", "ObjectID"})
  103. for i := 0; i < len(uploadObjectResult.ObjectResults); i++ {
  104. tb.AppendRow(table.Row{
  105. uploadObjectResult.ObjectResults[i].Info.Path,
  106. uploadObjectResult.ObjectResults[i].ObjectID,
  107. })
  108. }
  109. fmt.Print(tb.Render())
  110. return nil
  111. }
  112. if err != nil {
  113. return fmt.Errorf("wait uploading: %w", err)
  114. }
  115. }
  116. }
  117. func PackageUpdatePackage(ctx CommandContext, packageID cdssdk.PackageID, rootPath string) error {
  118. //userID := int64(0)
  119. var uploadFilePathes []string
  120. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  121. if err != nil {
  122. return nil
  123. }
  124. if !fi.IsDir() {
  125. uploadFilePathes = append(uploadFilePathes, fname)
  126. }
  127. return nil
  128. })
  129. if err != nil {
  130. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  131. }
  132. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  133. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartUpdatingPackage(0, packageID, objIter)
  134. if err != nil {
  135. return fmt.Errorf("update package %d failed, err: %w", packageID, err)
  136. }
  137. for {
  138. complete, _, err := ctx.Cmdline.Svc.PackageSvc().WaitUpdatingPackage(taskID, time.Second*5)
  139. if complete {
  140. if err != nil {
  141. return fmt.Errorf("updating package: %w", err)
  142. }
  143. return nil
  144. }
  145. if err != nil {
  146. return fmt.Errorf("wait updating: %w", err)
  147. }
  148. }
  149. }
  150. func PackageDeletePackage(ctx CommandContext, packageID cdssdk.PackageID) error {
  151. userID := cdssdk.UserID(0)
  152. err := ctx.Cmdline.Svc.PackageSvc().DeletePackage(userID, packageID)
  153. if err != nil {
  154. return fmt.Errorf("delete package %d failed, err: %w", packageID, err)
  155. }
  156. return nil
  157. }
  158. func PackageGetCachedNodes(ctx CommandContext, packageID cdssdk.PackageID, userID cdssdk.UserID) error {
  159. resp, err := ctx.Cmdline.Svc.PackageSvc().GetCachedNodes(userID, packageID)
  160. fmt.Printf("resp: %v\n", resp)
  161. if err != nil {
  162. return fmt.Errorf("get package %d cached nodes failed, err: %w", packageID, err)
  163. }
  164. return nil
  165. }
  166. func PackageGetLoadedNodes(ctx CommandContext, packageID cdssdk.PackageID, userID cdssdk.UserID) error {
  167. nodeIDs, err := ctx.Cmdline.Svc.PackageSvc().GetLoadedNodes(userID, packageID)
  168. fmt.Printf("nodeIDs: %v\n", nodeIDs)
  169. if err != nil {
  170. return fmt.Errorf("get package %d loaded nodes failed, err: %w", packageID, err)
  171. }
  172. return nil
  173. }
  174. func init() {
  175. commands.MustAdd(PackageListBucketPackages, "pkg", "ls")
  176. commands.MustAdd(PackageDownloadPackage, "pkg", "get")
  177. commands.MustAdd(PackageCreatePackage, "pkg", "new")
  178. commands.MustAdd(PackageUpdatePackage, "pkg", "update")
  179. commands.MustAdd(PackageDeletePackage, "pkg", "delete")
  180. commands.MustAdd(PackageGetCachedNodes, "pkg", "cached")
  181. commands.MustAdd(PackageGetLoadedNodes, "pkg", "loaded")
  182. }

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