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 8.7 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
2 years ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331
  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. stgsdk "gitlink.org.cn/cloudream/common/sdks/storage"
  10. "gitlink.org.cn/cloudream/storage/client/internal/config"
  11. "gitlink.org.cn/cloudream/storage/common/pkgs/iterator"
  12. )
  13. func PackageListBucketPackages(ctx CommandContext, bucketID int64) error {
  14. userID := int64(0)
  15. packages, err := ctx.Cmdline.Svc.BucketSvc().GetBucketPackages(userID, bucketID)
  16. if err != nil {
  17. return err
  18. }
  19. fmt.Printf("Find %d packages in bucket %d for user %d:\n", len(packages), bucketID, userID)
  20. tb := table.NewWriter()
  21. tb.AppendHeader(table.Row{"ID", "Name", "BucketID", "State", "Redundancy"})
  22. for _, obj := range packages {
  23. tb.AppendRow(table.Row{obj.PackageID, obj.Name, obj.BucketID, obj.State, obj.Redundancy})
  24. }
  25. fmt.Print(tb.Render())
  26. return nil
  27. }
  28. func PackageDownloadPackage(ctx CommandContext, outputDir string, packageID int64) error {
  29. err := os.MkdirAll(outputDir, os.ModePerm)
  30. if err != nil {
  31. return fmt.Errorf("create output directory %s failed, err: %w", outputDir, err)
  32. }
  33. // 下载文件
  34. objIter, err := ctx.Cmdline.Svc.PackageSvc().DownloadPackage(0, packageID)
  35. if err != nil {
  36. return fmt.Errorf("download object failed, err: %w", err)
  37. }
  38. defer objIter.Close()
  39. for {
  40. objInfo, err := objIter.MoveNext()
  41. if err == iterator.ErrNoMoreItem {
  42. break
  43. }
  44. if err != nil {
  45. return err
  46. }
  47. err = func() error {
  48. defer objInfo.File.Close()
  49. fullPath := filepath.Join(outputDir, objInfo.Object.Path)
  50. dirPath := filepath.Dir(fullPath)
  51. if err := os.MkdirAll(dirPath, 0755); err != nil {
  52. return fmt.Errorf("creating object dir: %w", err)
  53. }
  54. outputFile, err := os.Create(fullPath)
  55. if err != nil {
  56. return fmt.Errorf("creating object file: %w", err)
  57. }
  58. defer outputFile.Close()
  59. _, err = io.Copy(outputFile, objInfo.File)
  60. if err != nil {
  61. return fmt.Errorf("copy object data to local file failed, err: %w", err)
  62. }
  63. return nil
  64. }()
  65. if err != nil {
  66. return err
  67. }
  68. }
  69. return nil
  70. }
  71. func PackageUploadRepPackage(ctx CommandContext, rootPath string, bucketID int64, name string, repCount int, nodeAffinity []int64) error {
  72. rootPath = filepath.Clean(rootPath)
  73. var uploadFilePathes []string
  74. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  75. if err != nil {
  76. return nil
  77. }
  78. if !fi.IsDir() {
  79. uploadFilePathes = append(uploadFilePathes, fname)
  80. }
  81. return nil
  82. })
  83. if err != nil {
  84. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  85. }
  86. var nodeAff *int64
  87. if len(nodeAffinity) > 0 {
  88. nodeAff = &nodeAffinity[0]
  89. }
  90. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  91. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartCreatingRepPackage(0, bucketID, name, objIter, stgsdk.NewRepRedundancyInfo(repCount), 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().WaitCreatingRepPackage(taskID, time.Second*5)
  97. if complete {
  98. if err != nil {
  99. return fmt.Errorf("uploading rep object: %w", err)
  100. }
  101. tb := table.NewWriter()
  102. tb.AppendHeader(table.Row{"Path", "ObjectID", "FileHash"})
  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. uploadObjectResult.ObjectResults[i].FileHash,
  108. })
  109. }
  110. fmt.Print(tb.Render())
  111. return nil
  112. }
  113. if err != nil {
  114. return fmt.Errorf("wait uploading: %w", err)
  115. }
  116. }
  117. }
  118. func PackageUpdateRepPackage(ctx CommandContext, packageID int64, rootPath string) error {
  119. //userID := int64(0)
  120. var uploadFilePathes []string
  121. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  122. if err != nil {
  123. return nil
  124. }
  125. if !fi.IsDir() {
  126. uploadFilePathes = append(uploadFilePathes, fname)
  127. }
  128. return nil
  129. })
  130. if err != nil {
  131. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  132. }
  133. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  134. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartUpdatingRepPackage(0, packageID, objIter)
  135. if err != nil {
  136. return fmt.Errorf("update object %d failed, err: %w", packageID, err)
  137. }
  138. for {
  139. complete, _, err := ctx.Cmdline.Svc.PackageSvc().WaitUpdatingRepPackage(taskID, time.Second*5)
  140. if complete {
  141. if err != nil {
  142. return fmt.Errorf("updating rep object: %w", err)
  143. }
  144. return nil
  145. }
  146. if err != nil {
  147. return fmt.Errorf("wait updating: %w", err)
  148. }
  149. }
  150. }
  151. func PackageUploadECPackage(ctx CommandContext, rootPath string, bucketID int64, name string, ecName string, nodeAffinity []int64) error {
  152. var uploadFilePathes []string
  153. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  154. if err != nil {
  155. return nil
  156. }
  157. if !fi.IsDir() {
  158. uploadFilePathes = append(uploadFilePathes, fname)
  159. }
  160. return nil
  161. })
  162. if err != nil {
  163. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  164. }
  165. var nodeAff *int64
  166. if len(nodeAffinity) > 0 {
  167. nodeAff = &nodeAffinity[0]
  168. }
  169. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  170. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartCreatingECPackage(0, bucketID, name, objIter, stgsdk.NewECRedundancyInfo(ecName, config.Cfg().ECPacketSize), nodeAff)
  171. if err != nil {
  172. return fmt.Errorf("upload file data failed, err: %w", err)
  173. }
  174. for {
  175. complete, uploadObjectResult, err := ctx.Cmdline.Svc.PackageSvc().WaitCreatingRepPackage(taskID, time.Second*5)
  176. if complete {
  177. if err != nil {
  178. return fmt.Errorf("uploading ec package: %w", err)
  179. }
  180. tb := table.NewWriter()
  181. tb.AppendHeader(table.Row{"Path", "ObjectID", "FileHash"})
  182. for i := 0; i < len(uploadObjectResult.ObjectResults); i++ {
  183. tb.AppendRow(table.Row{
  184. uploadObjectResult.ObjectResults[i].Info.Path,
  185. uploadObjectResult.ObjectResults[i].ObjectID,
  186. uploadObjectResult.ObjectResults[i].FileHash,
  187. })
  188. }
  189. fmt.Print(tb.Render())
  190. return nil
  191. }
  192. if err != nil {
  193. return fmt.Errorf("wait uploading: %w", err)
  194. }
  195. }
  196. }
  197. func PackageUpdateECPackage(ctx CommandContext, packageID int64, rootPath string) error {
  198. //userID := int64(0)
  199. var uploadFilePathes []string
  200. err := filepath.WalkDir(rootPath, func(fname string, fi os.DirEntry, err error) error {
  201. if err != nil {
  202. return nil
  203. }
  204. if !fi.IsDir() {
  205. uploadFilePathes = append(uploadFilePathes, fname)
  206. }
  207. return nil
  208. })
  209. if err != nil {
  210. return fmt.Errorf("open directory %s failed, err: %w", rootPath, err)
  211. }
  212. objIter := iterator.NewUploadingObjectIterator(rootPath, uploadFilePathes)
  213. taskID, err := ctx.Cmdline.Svc.PackageSvc().StartUpdatingECPackage(0, packageID, objIter)
  214. if err != nil {
  215. return fmt.Errorf("update package %d failed, err: %w", packageID, err)
  216. }
  217. for {
  218. complete, _, err := ctx.Cmdline.Svc.PackageSvc().WaitUpdatingECPackage(taskID, time.Second*5)
  219. if complete {
  220. if err != nil {
  221. return fmt.Errorf("updating ec package: %w", err)
  222. }
  223. return nil
  224. }
  225. if err != nil {
  226. return fmt.Errorf("wait updating: %w", err)
  227. }
  228. }
  229. }
  230. func PackageDeletePackage(ctx CommandContext, packageID int64) error {
  231. userID := int64(0)
  232. err := ctx.Cmdline.Svc.PackageSvc().DeletePackage(userID, packageID)
  233. if err != nil {
  234. return fmt.Errorf("delete package %d failed, err: %w", packageID, err)
  235. }
  236. return nil
  237. }
  238. func PackageGetCachedNodes(ctx CommandContext, packageID int64, userID int64) error {
  239. resp, err := ctx.Cmdline.Svc.PackageSvc().GetCachedNodes(userID, packageID)
  240. fmt.Printf("resp: %v\n", resp)
  241. if err != nil {
  242. return fmt.Errorf("get package %d cached nodes failed, err: %w", packageID, err)
  243. }
  244. return nil
  245. }
  246. func PackageGetLoadedNodes(ctx CommandContext, packageID int64, userID int64) error {
  247. nodeIDs, err := ctx.Cmdline.Svc.PackageSvc().GetLoadedNodes(userID, packageID)
  248. fmt.Printf("nodeIDs: %v\n", nodeIDs)
  249. if err != nil {
  250. return fmt.Errorf("get package %d loaded nodes failed, err: %w", packageID, err)
  251. }
  252. return nil
  253. }
  254. func init() {
  255. commands.MustAdd(PackageListBucketPackages, "pkg", "ls")
  256. commands.MustAdd(PackageDownloadPackage, "pkg", "get")
  257. commands.MustAdd(PackageUploadRepPackage, "pkg", "new", "rep")
  258. commands.MustAdd(PackageUpdateRepPackage, "pkg", "update", "rep")
  259. commands.MustAdd(PackageUploadRepPackage, "pkg", "new", "ec")
  260. commands.MustAdd(PackageUpdateRepPackage, "pkg", "update", "ec")
  261. commands.MustAdd(PackageDeletePackage, "pkg", "delete")
  262. commands.MustAdd(PackageGetCachedNodes, "pkg", "cached")
  263. commands.MustAdd(PackageGetLoadedNodes, "pkg", "loaded")
  264. }

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