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.

fuse.go 8.1 kB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303
  1. package vfs
  2. import (
  3. "context"
  4. "strings"
  5. "time"
  6. "gitlink.org.cn/cloudream/common/utils/lo2"
  7. "gitlink.org.cn/cloudream/jcs-pub/client/internal/db"
  8. "gitlink.org.cn/cloudream/jcs-pub/client/internal/mount/fuse"
  9. "gitlink.org.cn/cloudream/jcs-pub/client/internal/mount/vfs/cache"
  10. clitypes "gitlink.org.cn/cloudream/jcs-pub/client/types"
  11. "gorm.io/gorm"
  12. )
  13. type FuseNode interface {
  14. PathComps() []string
  15. }
  16. func child(vfs *Vfs, ctx context.Context, parent FuseNode, name string) (fuse.FsEntry, error) {
  17. parentPathComps := parent.PathComps()
  18. childPathComps := lo2.AppendNew(parentPathComps, name)
  19. ca := vfs.cache.Stat(childPathComps)
  20. if ca == nil {
  21. var ret fuse.FsEntry
  22. d := vfs.db
  23. err := d.DoTx(func(tx db.SQLContext) error {
  24. pkg, err := d.Package().GetByFullName(tx, childPathComps[0], childPathComps[1])
  25. if err != nil {
  26. if err != gorm.ErrRecordNotFound {
  27. return err
  28. }
  29. return nil
  30. }
  31. objPath := clitypes.JoinObjectPath(childPathComps[2:]...)
  32. obj, err := d.Object().GetByPath(tx, pkg.PackageID, objPath)
  33. if err == nil {
  34. ret = newFileFromObject(vfs, childPathComps, obj)
  35. return nil
  36. }
  37. if err != gorm.ErrRecordNotFound {
  38. return err
  39. }
  40. err = d.Object().HasObjectWithPrefix(tx, pkg.PackageID, objPath+clitypes.ObjectPathSeparator)
  41. if err == nil {
  42. dir := vfs.cache.LoadDir(childPathComps, &cache.CreateDirOption{
  43. ModTime: time.Now(),
  44. })
  45. if dir == nil {
  46. return nil
  47. }
  48. ret = newDirFromCache(dir.Info(), vfs)
  49. return nil
  50. }
  51. if err == gorm.ErrRecordNotFound {
  52. return nil
  53. }
  54. return err
  55. })
  56. if err != nil {
  57. return nil, err
  58. }
  59. if ret == nil {
  60. return nil, fuse.ErrNotExists
  61. }
  62. return ret, nil
  63. }
  64. if ca.IsDir {
  65. return newDirFromCache(*ca, vfs), nil
  66. }
  67. return newFileFromCache(*ca, vfs), nil
  68. }
  69. func listChildren(vfs *Vfs, ctx context.Context, parent FuseNode) ([]fuse.FsEntry, error) {
  70. var ens []fuse.FsEntry
  71. myPathComps := parent.PathComps()
  72. infos := vfs.cache.StatMany(myPathComps)
  73. dbEntries := make(map[string]fuse.FsEntry)
  74. d := vfs.db
  75. d.DoTx(func(tx db.SQLContext) error {
  76. pkg, err := d.Package().GetByFullName(tx, myPathComps[0], myPathComps[1])
  77. if err != nil {
  78. return err
  79. }
  80. objPath := clitypes.JoinObjectPath(myPathComps[2:]...)
  81. objPrefix := objPath
  82. if objPath != "" {
  83. objPrefix += clitypes.ObjectPathSeparator
  84. }
  85. objs, coms, err := d.Object().GetByPrefixGrouped(tx, pkg.PackageID, objPrefix)
  86. if err != nil {
  87. return err
  88. }
  89. for _, dir := range coms {
  90. dir = strings.TrimSuffix(dir, clitypes.ObjectPathSeparator)
  91. pathComps := lo2.AppendNew(myPathComps, clitypes.BaseName(dir))
  92. cd := vfs.cache.LoadDir(pathComps, &cache.CreateDirOption{
  93. ModTime: time.Now(),
  94. })
  95. if cd == nil {
  96. continue
  97. }
  98. dbEntries[dir] = newDirFromCache(cd.Info(), vfs)
  99. }
  100. for _, obj := range objs {
  101. pathComps := lo2.AppendNew(myPathComps, clitypes.BaseName(obj.Path))
  102. file := newFileFromObject(vfs, pathComps, obj)
  103. dbEntries[file.Name()] = file
  104. }
  105. return nil
  106. })
  107. for _, c := range infos {
  108. delete(dbEntries, c.PathComps[len(c.PathComps)-1])
  109. if c.IsDir {
  110. ens = append(ens, newDirFromCache(c, vfs))
  111. } else {
  112. ens = append(ens, newFileFromCache(c, vfs))
  113. }
  114. }
  115. for _, e := range dbEntries {
  116. ens = append(ens, e)
  117. }
  118. return ens, nil
  119. }
  120. func newDir(vfs *Vfs, ctx context.Context, name string, parent FuseNode) (fuse.FsDir, error) {
  121. cache := vfs.cache.CreateDir(lo2.AppendNew(parent.PathComps(), name))
  122. if cache == nil {
  123. return nil, fuse.ErrPermission
  124. }
  125. return newDirFromCache(cache.Info(), vfs), nil
  126. }
  127. func newFile(vfs *Vfs, ctx context.Context, name string, parent FuseNode, flags uint32) (fuse.FileHandle, uint32, error) {
  128. ch := vfs.cache.CreateFile(lo2.AppendNew(parent.PathComps(), name))
  129. if ch == nil {
  130. return nil, 0, fuse.ErrPermission
  131. }
  132. defer ch.Release()
  133. if !ch.LevelUp(cache.LevelComplete) {
  134. return nil, 0, fuse.ErrIOError
  135. }
  136. // Open之后会给cache的引用计数额外+1,即使cache先于FileHandle被关闭,
  137. // 也有有FileHandle的计数保持cache的有效性
  138. fileNode := newFileFromCache(ch.Info(), vfs)
  139. hd := ch.Open(flags)
  140. return newFileHandle(fileNode, hd), flags, nil
  141. }
  142. func removeChild(vfs *Vfs, ctx context.Context, name string, parent FuseNode) error {
  143. pathComps := lo2.AppendNew(parent.PathComps(), name)
  144. joinedPath := clitypes.JoinObjectPath(pathComps[2:]...)
  145. d := vfs.db
  146. // TODO 生成系统事件
  147. return vfs.db.DoTx(func(tx db.SQLContext) error {
  148. pkg, err := d.Package().GetByFullName(tx, pathComps[0], pathComps[1])
  149. if err == nil {
  150. err := d.Object().HasObjectWithPrefix(tx, pkg.PackageID, joinedPath+clitypes.ObjectPathSeparator)
  151. if err == nil {
  152. return fuse.ErrNotEmpty
  153. }
  154. if err != gorm.ErrRecordNotFound {
  155. return err
  156. }
  157. // 存储系统不会保存目录结构,所以这里是尝试删除同名文件
  158. err = d.Object().DeleteCompleteByPath(tx, pkg.PackageID, joinedPath)
  159. if err != nil && err != gorm.ErrRecordNotFound {
  160. return err
  161. }
  162. } else if err != gorm.ErrRecordNotFound {
  163. return err
  164. }
  165. return vfs.cache.Remove(pathComps)
  166. })
  167. }
  168. func moveChild(vfs *Vfs, ctx context.Context, oldName string, oldParent FuseNode, newName string, newParent FuseNode) error {
  169. newParentPath := newParent.PathComps()
  170. newChildPath := lo2.AppendNew(newParentPath, newName)
  171. newChildPathJoined := clitypes.JoinObjectPath(newChildPath[2:]...)
  172. // 不允许移动任何内容到Package层级以上
  173. if len(newParentPath) < 2 {
  174. return fuse.ErrNotSupported
  175. }
  176. oldChildPath := lo2.AppendNew(oldParent.PathComps(), oldName)
  177. oldChildPathJoined := clitypes.JoinObjectPath(oldChildPath[2:]...)
  178. // 先更新远程,再更新本地,因为远程使用事务更新,可以回滚,而本地不行
  179. return vfs.db.DoTx(func(tx db.SQLContext) error {
  180. err := moveRemote(vfs, tx, oldChildPath, newParentPath, oldChildPathJoined, newChildPathJoined)
  181. if err == fuse.ErrExists {
  182. return err
  183. }
  184. if err != nil && err != fuse.ErrNotExists {
  185. return err
  186. }
  187. err2 := vfs.cache.Move(oldChildPath, newChildPath)
  188. if err2 == fuse.ErrNotExists {
  189. if err == fuse.ErrNotExists {
  190. return fuse.ErrNotExists
  191. }
  192. return nil
  193. }
  194. return err2
  195. })
  196. }
  197. func moveRemote(vfs *Vfs, tx db.SQLContext, oldChildPath []string, newParentPath []string, oldChildPathJoined string, newChildPathJoined string) error {
  198. d := vfs.db
  199. newPkg, err := d.Package().GetByFullName(tx, newParentPath[0], newParentPath[1])
  200. if err != nil {
  201. if err == gorm.ErrRecordNotFound {
  202. return fuse.ErrNotExists
  203. }
  204. return err
  205. }
  206. // 检查目的文件或文件夹是否已经存在
  207. _, err = d.Object().GetByPath(tx, newPkg.PackageID, newChildPathJoined)
  208. if err == nil {
  209. return fuse.ErrExists
  210. }
  211. err = d.Object().HasObjectWithPrefix(tx, newPkg.PackageID, newChildPathJoined+clitypes.ObjectPathSeparator)
  212. if err == nil {
  213. return fuse.ErrExists
  214. }
  215. if err != gorm.ErrRecordNotFound {
  216. return err
  217. }
  218. // 按理来说还需要检查远程文件所在的文件夹是否存在,但对象存储是不存文件夹的,所以不检查,导致的后果就是移动时会创建不存在的文件夹
  219. oldPkg, err := d.Package().GetByFullName(tx, oldChildPath[0], oldChildPath[1])
  220. if err != nil {
  221. if err == gorm.ErrRecordNotFound {
  222. return fuse.ErrNotExists
  223. }
  224. return err
  225. }
  226. // 都不存在,就开始移动文件
  227. oldObj, err := d.Object().GetByPath(tx, oldPkg.PackageID, oldChildPathJoined)
  228. if err == nil {
  229. oldObj.PackageID = newPkg.PackageID
  230. oldObj.Path = newChildPathJoined
  231. return d.Object().BatchUpdate(tx, []clitypes.Object{oldObj})
  232. }
  233. if err != gorm.ErrRecordNotFound {
  234. return err
  235. }
  236. err = d.Object().HasObjectWithPrefix(tx, oldPkg.PackageID, oldChildPathJoined+clitypes.ObjectPathSeparator)
  237. if err == nil {
  238. return d.Object().MoveByPrefix(tx,
  239. oldPkg.PackageID, oldChildPathJoined+clitypes.ObjectPathSeparator,
  240. newPkg.PackageID, newChildPathJoined+clitypes.ObjectPathSeparator,
  241. )
  242. }
  243. if err == gorm.ErrRecordNotFound {
  244. return fuse.ErrNotExists
  245. }
  246. return err
  247. }

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