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

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309
  1. package cmdline
  2. import (
  3. "fmt"
  4. "io"
  5. "os"
  6. "path/filepath"
  7. "strings"
  8. "time"
  9. "github.com/jedib0t/go-pretty/v6/table"
  10. "github.com/juju/ratelimit"
  11. "gitlink.org.cn/cloudream/client/internal/task"
  12. myio "gitlink.org.cn/cloudream/common/utils/io"
  13. )
  14. func ObjectListBucketObjects(ctx CommandContext, bucketID int) error {
  15. userID := 0
  16. objects, err := ctx.Cmdline.Svc.BucketSvc().GetBucketObjects(userID, bucketID)
  17. if err != nil {
  18. return err
  19. }
  20. fmt.Printf("Find %d objects in bucket %d for user %d:\n", len(objects), bucketID, userID)
  21. tb := table.NewWriter()
  22. tb.AppendHeader(table.Row{"ID", "Name", "Size", "BucketID", "State", "Redundancy"})
  23. for _, obj := range objects {
  24. tb.AppendRow(table.Row{obj.ObjectID, obj.Name, obj.BucketID, obj.State, obj.FileSize, obj.Redundancy})
  25. }
  26. fmt.Print(tb.Render())
  27. return nil
  28. }
  29. func ObjectDownloadObject(ctx CommandContext, localFilePath string, objectID int) error {
  30. // 创建本地文件
  31. curExecPath, err := os.Executable()
  32. if err != nil {
  33. return fmt.Errorf("get executable directory failed, err: %w", err)
  34. }
  35. outputFilePath := filepath.Join(filepath.Dir(curExecPath), localFilePath)
  36. outputFileDir := filepath.Dir(outputFilePath)
  37. err = os.MkdirAll(outputFileDir, os.ModePerm)
  38. if err != nil {
  39. return fmt.Errorf("create output file directory %s failed, err: %w", outputFileDir, err)
  40. }
  41. outputFile, err := os.Create(outputFilePath)
  42. if err != nil {
  43. return fmt.Errorf("create output file %s failed, err: %w", outputFilePath, err)
  44. }
  45. defer outputFile.Close()
  46. // 下载文件
  47. reader, err := ctx.Cmdline.Svc.ObjectSvc().DownloadObject(0, objectID)
  48. if err != nil {
  49. return fmt.Errorf("download object failed, err: %w", err)
  50. }
  51. defer reader.Close()
  52. bkt := ratelimit.NewBucketWithRate(10*1024, 10*1024)
  53. _, err = io.Copy(outputFile, ratelimit.Reader(reader, bkt))
  54. if err != nil {
  55. // TODO 写入到文件失败,是否要考虑删除这个不完整的文件?
  56. return fmt.Errorf("copy object data to local file failed, err: %w", err)
  57. }
  58. return nil
  59. }
  60. func ObjectDownloadObjectDir(ctx CommandContext, localFilePath string, dirName string) error {
  61. /* // 创建本地文件夹
  62. curExecPath, err := os.Executable()
  63. if err != nil {
  64. return fmt.Errorf("get executable directory failed, err: %w", err)
  65. }
  66. outputFilePath := filepath.Join(filepath.Dir(curExecPath), localFilePath)
  67. outputFileDir := filepath.Dir(outputFilePath)
  68. err = os.MkdirAll(outputFileDir, os.ModePerm)
  69. if err != nil {
  70. return fmt.Errorf("create output file directory %s failed, err: %w", outputFileDir, err)
  71. }
  72. outputFile, err := os.Create(outputFilePath)
  73. if err != nil {
  74. return fmt.Errorf("create output file %s failed, err: %w", outputFilePath, err)
  75. }
  76. defer outputFile.Close()
  77. // 下载文件
  78. reader, err := ctx.Cmdline.Svc.ObjectSvc().DownloadObject(0, objectID)
  79. if err != nil {
  80. return fmt.Errorf("download object failed, err: %w", err)
  81. }
  82. defer reader.Close()
  83. bkt := ratelimit.NewBucketWithRate(10*1024, 10*1024)
  84. _, err = io.Copy(outputFile, ratelimit.Reader(reader, bkt))
  85. if err != nil {
  86. // TODO 写入到文件失败,是否要考虑删除这个不完整的文件?
  87. return fmt.Errorf("copy object data to local file failed, err: %w", err)
  88. } */
  89. return nil
  90. }
  91. func ObjectUploadRepObject(ctx CommandContext, localFilePath string, bucketID int, objectName string, repCount int) error {
  92. file, err := os.Open(localFilePath)
  93. if err != nil {
  94. return fmt.Errorf("open file %s failed, err: %w", localFilePath, err)
  95. }
  96. defer file.Close()
  97. fileInfo, err := file.Stat()
  98. if err != nil {
  99. return fmt.Errorf("get file %s state failed, err: %w", localFilePath, err)
  100. }
  101. fileSize := fileInfo.Size()
  102. // TODO 测试用
  103. bkt := ratelimit.NewBucketWithRate(10*1024, 10*1024)
  104. uploadObject := task.UploadObject{
  105. ObjectName: objectName,
  106. File: myio.WithCloser(ratelimit.Reader(file, bkt),
  107. func(reader io.Reader) error {
  108. return file.Close()
  109. }),
  110. FileSize: fileSize,
  111. }
  112. uploadObjects := []task.UploadObject{uploadObject}
  113. taskID, err := ctx.Cmdline.Svc.ObjectSvc().StartUploadingRepObjects(0, bucketID, uploadObjects, repCount)
  114. if err != nil {
  115. return fmt.Errorf("upload file data failed, err: %w", err)
  116. }
  117. for {
  118. complete, UploadObjectResult, err := ctx.Cmdline.Svc.ObjectSvc().WaitUploadingRepObjects(taskID, time.Second*5)
  119. if complete {
  120. if err != nil {
  121. return fmt.Errorf("uploading rep object: %w", err)
  122. }
  123. fmt.Print(UploadObjectResult.UploadRepResults[0].ResultFileHash)
  124. return nil
  125. }
  126. if err != nil {
  127. return fmt.Errorf("wait uploading: %w", err)
  128. }
  129. }
  130. }
  131. func ObjectUploadRepObjectDir(ctx CommandContext, localDirPath string, bucketID int, repCount int) error {
  132. var uploadFiles []task.UploadObject
  133. var uploadFile task.UploadObject
  134. err := filepath.Walk(localDirPath, func(fname string, fi os.FileInfo, err error) error {
  135. if !fi.IsDir() {
  136. file, err := os.Open(fname)
  137. if err != nil {
  138. return fmt.Errorf("open file %s failed, err: %w", fname, err)
  139. }
  140. // TODO 测试用
  141. bkt := ratelimit.NewBucketWithRate(10*1024, 10*1024)
  142. uploadFile = task.UploadObject{
  143. ObjectName: strings.Replace(fname, "\\", "/", -1),
  144. File: myio.WithCloser(ratelimit.Reader(file, bkt),
  145. func(reader io.Reader) error {
  146. return file.Close()
  147. }),
  148. FileSize: fi.Size(),
  149. }
  150. uploadFiles = append(uploadFiles, uploadFile)
  151. }
  152. return nil
  153. })
  154. if err != nil {
  155. return fmt.Errorf("open directory %s failed, err: %w", localDirPath, err)
  156. }
  157. // 遍历 关闭文件流
  158. defer func() {
  159. for _, uploadFile := range uploadFiles {
  160. uploadFile.File.Close()
  161. }
  162. }()
  163. taskID, err := ctx.Cmdline.Svc.ObjectSvc().StartUploadingRepObjects(0, bucketID, uploadFiles, repCount)
  164. if err != nil {
  165. return fmt.Errorf("upload file data failed, err: %w", err)
  166. }
  167. for {
  168. complete, UploadObjectResult, err := ctx.Cmdline.Svc.ObjectSvc().WaitUploadingRepObjects(taskID, time.Second*5)
  169. if complete {
  170. if err != nil {
  171. return fmt.Errorf("uploading rep object: %w", err)
  172. }
  173. tb := table.NewWriter()
  174. if UploadObjectResult.IsUploading {
  175. tb.AppendHeader(table.Row{"ObjectID", "ObjectName", "FileHash"})
  176. for i := 0; i < len(UploadObjectResult.UploadObjects); i++ {
  177. tb.AppendRow(table.Row{UploadObjectResult.UploadRepResults[i].ObjectID, UploadObjectResult.UploadObjects[i].ObjectName, UploadObjectResult.UploadRepResults[i].ResultFileHash})
  178. }
  179. fmt.Print(tb.Render())
  180. } else {
  181. fmt.Println("The folder upload failed. Some files do not meet the upload requirements.")
  182. tb.AppendHeader(table.Row{"ObjectName", "Error"})
  183. for i := 0; i < len(UploadObjectResult.UploadObjects); i++ {
  184. if UploadObjectResult.UploadRepResults[i].Error != nil {
  185. tb.AppendRow(table.Row{UploadObjectResult.UploadObjects[i].ObjectName, UploadObjectResult.UploadRepResults[i].Error})
  186. }
  187. }
  188. fmt.Print(tb.Render())
  189. }
  190. return nil
  191. }
  192. if err != nil {
  193. return fmt.Errorf("wait uploading: %w", err)
  194. }
  195. }
  196. }
  197. func ObjectEcWrite(ctx CommandContext, localFilePath string, bucketID int, objectName string, ecName string) error {
  198. // TODO
  199. panic("not implement yet")
  200. }
  201. func ObjectUpdateRepObject(ctx CommandContext, objectID int, filePath string) error {
  202. userID := 0
  203. file, err := os.Open(filePath)
  204. if err != nil {
  205. return fmt.Errorf("open file %s failed, err: %w", filePath, err)
  206. }
  207. defer file.Close()
  208. fileInfo, err := file.Stat()
  209. if err != nil {
  210. return fmt.Errorf("get file %s state failed, err: %w", filePath, err)
  211. }
  212. fileSize := fileInfo.Size()
  213. // TODO 测试用
  214. bkt := ratelimit.NewBucketWithRate(10*1024, 10*1024)
  215. taskID, err := ctx.Cmdline.Svc.ObjectSvc().StartUpdatingRepObject(userID, objectID,
  216. myio.WithCloser(ratelimit.Reader(file, bkt),
  217. func(reader io.Reader) error {
  218. return file.Close()
  219. }), fileSize)
  220. if err != nil {
  221. return fmt.Errorf("update object %d failed, err: %w", objectID, err)
  222. }
  223. for {
  224. complete, err := ctx.Cmdline.Svc.ObjectSvc().WaitUpdatingRepObject(taskID, time.Second*5)
  225. if complete {
  226. if err != nil {
  227. return fmt.Errorf("updating rep object: %w", err)
  228. }
  229. return nil
  230. }
  231. if err != nil {
  232. return fmt.Errorf("wait updating: %w", err)
  233. }
  234. }
  235. }
  236. func ObjectDeleteObject(ctx CommandContext, objectID int) error {
  237. userID := 0
  238. err := ctx.Cmdline.Svc.ObjectSvc().DeleteObject(userID, objectID)
  239. if err != nil {
  240. return fmt.Errorf("delete object %d failed, err: %w", objectID, err)
  241. }
  242. return nil
  243. }
  244. func init() {
  245. commands.MustAdd(ObjectListBucketObjects, "object", "ls")
  246. commands.MustAdd(ObjectUploadRepObject, "object", "new", "rep")
  247. commands.MustAdd(ObjectUploadRepObjectDir, "object", "new", "dir")
  248. commands.MustAdd(ObjectDownloadObject, "object", "get")
  249. commands.MustAdd(ObjectDownloadObjectDir, "object", "get", "dir")
  250. commands.MustAdd(ObjectUpdateRepObject, "object", "update", "rep")
  251. commands.MustAdd(ObjectDeleteObject, "object", "delete")
  252. }

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