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

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  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. "gitlink.org.cn/cloudream/common/models"
  10. "gitlink.org.cn/cloudream/storage-common/pkgs/iterator"
  11. )
  12. func PackageListBucketPackages(ctx CommandContext, bucketID int64) error {
  13. userID := int64(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", "Redundancy"})
  21. for _, obj := range packages {
  22. tb.AppendRow(table.Row{obj.PackageID, obj.Name, obj.BucketID, obj.State, obj.Redundancy})
  23. }
  24. fmt.Print(tb.Render())
  25. return nil
  26. }
  27. func PackageDownloadPackage(ctx CommandContext, outputDir string, packageID int64) 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. defer objInfo.File.Close()
  47. fullPath := filepath.Join(outputDir, objInfo.Object.Path)
  48. dirPath := filepath.Dir(fullPath)
  49. if err := os.MkdirAll(dirPath, 0755); err != nil {
  50. return fmt.Errorf("creating object dir: %w", err)
  51. }
  52. outputFile, err := os.Create(fullPath)
  53. if err != nil {
  54. return fmt.Errorf("creating object file: %w", err)
  55. }
  56. defer outputFile.Close()
  57. _, err = io.Copy(outputFile, objInfo.File)
  58. if err != nil {
  59. return fmt.Errorf("copy object data to local file failed, err: %w", err)
  60. }
  61. }
  62. return nil
  63. }
  64. func PackageUploadRepPackage(ctx CommandContext, rootPath string, bucketID int64, name string, repCount int) error {
  65. rootPath = filepath.Clean(rootPath)
  66. var uploadFilePathes []string
  67. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  68. if err != nil {
  69. return nil
  70. }
  71. if !fi.IsDir() {
  72. uploadFilePathes = append(uploadFilePathes, fname)
  73. }
  74. return nil
  75. })
  76. if err != nil {
  77. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  78. }
  79. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  80. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartCreatingRepPackage(0, bucketID, name, objIter, models.NewRepRedundancyInfo(repCount))
  81. if err != nil {
  82. return fmt.Errorf("upload file data failed, err: %w", err)
  83. }
  84. for {
  85. complete, uploadObjectResult, err := ctx.Cmdline.Svc.PackageSvc().WaitCreatingRepPackage(taskID, time.Second*5)
  86. if complete {
  87. if err != nil {
  88. return fmt.Errorf("uploading rep object: %w", err)
  89. }
  90. tb := table.NewWriter()
  91. tb.AppendHeader(table.Row{"Path", "ObjectID", "FileHash"})
  92. for i := 0; i < len(uploadObjectResult.ObjectResults); i++ {
  93. tb.AppendRow(table.Row{
  94. uploadObjectResult.ObjectResults[i].Info.Path,
  95. uploadObjectResult.ObjectResults[i].ObjectID,
  96. uploadObjectResult.ObjectResults[i].FileHash,
  97. })
  98. }
  99. fmt.Print(tb.Render())
  100. return nil
  101. }
  102. if err != nil {
  103. return fmt.Errorf("wait uploading: %w", err)
  104. }
  105. }
  106. }
  107. func PackageUpdateRepPackage(ctx CommandContext, packageID int64, rootPath string) error {
  108. //userID := int64(0)
  109. var uploadFilePathes []string
  110. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  111. if err != nil {
  112. return nil
  113. }
  114. if !fi.IsDir() {
  115. uploadFilePathes = append(uploadFilePathes, fname)
  116. }
  117. return nil
  118. })
  119. if err != nil {
  120. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  121. }
  122. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  123. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartUpdatingRepPackage(0, packageID, objIter)
  124. if err != nil {
  125. return fmt.Errorf("update object %d failed, err: %w", packageID, err)
  126. }
  127. for {
  128. complete, _, err := ctx.Cmdline.Svc.PackageSvc().WaitUpdatingRepPackage(taskID, time.Second*5)
  129. if complete {
  130. if err != nil {
  131. return fmt.Errorf("updating rep object: %w", err)
  132. }
  133. return nil
  134. }
  135. if err != nil {
  136. return fmt.Errorf("wait updating: %w", err)
  137. }
  138. }
  139. }
  140. func PackageUploadECPackage(ctx CommandContext, rootPath string, bucketID int64, name string, ecName string) error {
  141. var uploadFilePathes []string
  142. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  143. if err != nil {
  144. return nil
  145. }
  146. if !fi.IsDir() {
  147. uploadFilePathes = append(uploadFilePathes, fname)
  148. }
  149. return nil
  150. })
  151. if err != nil {
  152. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  153. }
  154. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  155. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartCreatingECPackage(0, bucketID, name, objIter, models.NewECRedundancyInfo(ecName))
  156. if err != nil {
  157. return fmt.Errorf("upload file data failed, err: %w", err)
  158. }
  159. for {
  160. complete, uploadObjectResult, err := ctx.Cmdline.Svc.PackageSvc().WaitCreatingRepPackage(taskID, time.Second*5)
  161. if complete {
  162. if err != nil {
  163. return fmt.Errorf("uploading ec package: %w", err)
  164. }
  165. tb := table.NewWriter()
  166. tb.AppendHeader(table.Row{"Path", "ObjectID", "FileHash"})
  167. for i := 0; i < len(uploadObjectResult.ObjectResults); i++ {
  168. tb.AppendRow(table.Row{
  169. uploadObjectResult.ObjectResults[i].Info.Path,
  170. uploadObjectResult.ObjectResults[i].ObjectID,
  171. uploadObjectResult.ObjectResults[i].FileHash,
  172. })
  173. }
  174. fmt.Print(tb.Render())
  175. return nil
  176. }
  177. if err != nil {
  178. return fmt.Errorf("wait uploading: %w", err)
  179. }
  180. }
  181. }
  182. func PackageUpdateECPackage(ctx CommandContext, packageID int64, rootPath string) error {
  183. //userID := int64(0)
  184. var uploadFilePathes []string
  185. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  186. if err != nil {
  187. return nil
  188. }
  189. if !fi.IsDir() {
  190. uploadFilePathes = append(uploadFilePathes, fname)
  191. }
  192. return nil
  193. })
  194. if err != nil {
  195. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  196. }
  197. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  198. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartUpdatingECPackage(0, packageID, objIter)
  199. if err != nil {
  200. return fmt.Errorf("update package %d failed, err: %w", packageID, err)
  201. }
  202. for {
  203. complete, _, err := ctx.Cmdline.Svc.PackageSvc().WaitUpdatingECPackage(taskID, time.Second*5)
  204. if complete {
  205. if err != nil {
  206. return fmt.Errorf("updating ec package: %w", err)
  207. }
  208. return nil
  209. }
  210. if err != nil {
  211. return fmt.Errorf("wait updating: %w", err)
  212. }
  213. }
  214. }
  215. func PackageDeletePackage(ctx CommandContext, packageID int64) error {
  216. userID := int64(0)
  217. err := ctx.Cmdline.Svc.PackageSvc().DeletePackage(userID, packageID)
  218. if err != nil {
  219. return fmt.Errorf("delete package %d failed, err: %w", packageID, err)
  220. }
  221. return nil
  222. }
  223. func init() {
  224. commands.MustAdd(PackageListBucketPackages, "pkg", "ls")
  225. commands.MustAdd(PackageDownloadPackage, "pkg", "get")
  226. commands.MustAdd(PackageUploadRepPackage, "pkg", "new", "rep")
  227. commands.MustAdd(PackageUpdateRepPackage, "pkg", "update", "rep")
  228. commands.MustAdd(PackageUploadRepPackage, "pkg", "new", "ec")
  229. commands.MustAdd(PackageUpdateRepPackage, "pkg", "update", "ec")
  230. commands.MustAdd(PackageDeletePackage, "pkg", "delete")
  231. }

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