Browse Source

Merge branch 'master' into feature_gxh

pull/1/head
Sydonian 2 years ago
parent
commit
322e11c5ee
74 changed files with 72 additions and 439 deletions
  1. +0
    -26
      consts/consts.go
  2. +1
    -4
      go.mod
  3. +0
    -12
      go.sum
  4. +9
    -3
      magefiles/common.go
  5. +5
    -0
      magefiles/targets/targets.go
  6. +0
    -37
      models/models.go
  7. +1
    -1
      pkgs/actor/actor.go
  8. +0
    -0
      pkgs/actor/actor_test.go
  9. +1
    -1
      pkgs/cmdtrie/command_trie.go
  10. +0
    -0
      pkgs/cmdtrie/command_trie_test.go
  11. +0
    -0
      pkgs/distlock/config.go
  12. +0
    -0
      pkgs/distlock/const.go
  13. +0
    -0
      pkgs/distlock/distlock.go
  14. +1
    -1
      pkgs/distlock/lockprovider/ipfs_lock.go
  15. +1
    -1
      pkgs/distlock/lockprovider/ipfs_lock_test.go
  16. +1
    -1
      pkgs/distlock/lockprovider/lock_compatibility_table.go
  17. +1
    -1
      pkgs/distlock/lockprovider/lock_compatibility_table_test.go
  18. +1
    -1
      pkgs/distlock/lockprovider/metadata_lock.go
  19. +1
    -1
      pkgs/distlock/lockprovider/storage_lock.go
  20. +0
    -0
      pkgs/distlock/lockprovider/string_lock_target.go
  21. +0
    -0
      pkgs/distlock/lockprovider/string_lock_target_test.go
  22. +2
    -2
      pkgs/distlock/reqbuilder/ipfs.go
  23. +2
    -2
      pkgs/distlock/reqbuilder/lock_request_builder.go
  24. +1
    -1
      pkgs/distlock/reqbuilder/metadata.go
  25. +2
    -2
      pkgs/distlock/reqbuilder/metadata_bucket.go
  26. +2
    -2
      pkgs/distlock/reqbuilder/metadata_cache.go
  27. +2
    -2
      pkgs/distlock/reqbuilder/metadata_node.go
  28. +2
    -2
      pkgs/distlock/reqbuilder/metadata_object.go
  29. +2
    -2
      pkgs/distlock/reqbuilder/metadata_object_block.go
  30. +2
    -2
      pkgs/distlock/reqbuilder/metadata_object_rep.go
  31. +2
    -2
      pkgs/distlock/reqbuilder/metadata_storage_object.go
  32. +2
    -2
      pkgs/distlock/reqbuilder/metadata_user_bucket.go
  33. +2
    -2
      pkgs/distlock/reqbuilder/metadata_user_storage.go
  34. +2
    -2
      pkgs/distlock/reqbuilder/storage.go
  35. +4
    -4
      pkgs/distlock/service/init_providers.go
  36. +2
    -2
      pkgs/distlock/service/internal/lease_actor.go
  37. +2
    -2
      pkgs/distlock/service/internal/main_actor.go
  38. +4
    -4
      pkgs/distlock/service/internal/providers_actor.go
  39. +4
    -4
      pkgs/distlock/service/internal/retry_actor.go
  40. +0
    -0
      pkgs/distlock/service/internal/utils.go
  41. +0
    -0
      pkgs/distlock/service/internal/utils_test.go
  42. +1
    -1
      pkgs/distlock/service/internal/watch_etcd_actor.go
  43. +1
    -1
      pkgs/distlock/service/mutex.go
  44. +3
    -3
      pkgs/distlock/service/service.go
  45. +0
    -0
      pkgs/event/event.go
  46. +0
    -0
      pkgs/event/executor.go
  47. +0
    -0
      pkgs/future/future.go
  48. +0
    -0
      pkgs/future/future_test.go
  49. +0
    -0
      pkgs/future/set_value_future.go
  50. +0
    -0
      pkgs/future/set_void_future.go
  51. +0
    -0
      pkgs/logger/config.go
  52. +0
    -0
      pkgs/logger/global_logger.go
  53. +0
    -0
      pkgs/logger/logger.go
  54. +0
    -0
      pkgs/logger/logger_test.go
  55. +0
    -0
      pkgs/logger/logrus_logger.go
  56. +0
    -0
      pkgs/logger/utils.go
  57. +2
    -2
      pkgs/mq/client.go
  58. +0
    -0
      pkgs/mq/message.go
  59. +0
    -0
      pkgs/mq/message_dispatcher.go
  60. +0
    -0
      pkgs/mq/message_test.go
  61. +0
    -0
      pkgs/mq/response.go
  62. +0
    -0
      pkgs/mq/server.go
  63. +0
    -0
      pkgs/task/manager.go
  64. +0
    -0
      pkgs/task/task.go
  65. +0
    -0
      pkgs/tickevent/executor.go
  66. +0
    -0
      pkgs/tickevent/tick_event.go
  67. +0
    -0
      pkgs/trie/trie.go
  68. +0
    -0
      pkgs/trie/trie_test.go
  69. +0
    -0
      pkgs/typedispatcher/type_dispatcher.go
  70. +0
    -74
      utils/config.go
  71. +1
    -1
      utils/config/config.go
  72. +0
    -123
      utils/grpc/file_transport.go
  73. +0
    -82
      utils/ping.go
  74. +0
    -21
      utils/utils.go

+ 0
- 26
consts/consts.go View File

@@ -1,26 +0,0 @@
package consts

const (
IPFSStateOK = "OK"

StorageDirectoryStateOK = "OK"

NodeStateNormal = "Normal"
NodeStateUnavailable = "Unavailable"
)

const (
ObjectStateNormal = "Normal"
ObjectStateDeleted = "Deleted"
)

const (
StorageObjectStateNormal = "Normal"
StorageObjectStateDeleted = "Deleted"
StorageObjectStateOutdated = "Outdated"
)

const (
CacheStatePinned = "Pinned"
CacheStateTemp = "Temp"
)

+ 1
- 4
go.mod View File

@@ -1,11 +1,9 @@
module gitlink.org.cn/cloudream/common

go 1.18
go 1.20

require (
github.com/antonfisher/nested-logrus-formatter v1.3.1
github.com/beevik/etree v1.2.0
github.com/go-ping/ping v1.1.0
github.com/google/uuid v1.3.0
github.com/hashicorp/go-multierror v1.1.1
github.com/imdario/mergo v0.3.15
@@ -63,7 +61,6 @@ require (
go.uber.org/zap v1.24.0 // indirect
golang.org/x/crypto v0.6.0 // indirect
golang.org/x/net v0.8.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/sys v0.6.0 // indirect
golang.org/x/text v0.8.0 // indirect
google.golang.org/genproto v0.0.0-20230403163135-c38d8f061ccd // indirect


+ 0
- 12
go.sum View File

@@ -1,7 +1,5 @@
github.com/antonfisher/nested-logrus-formatter v1.3.1 h1:NFJIr+pzwv5QLHTPyKz9UMEoHck02Q9L0FP13b/xSbQ=
github.com/antonfisher/nested-logrus-formatter v1.3.1/go.mod h1:6WTfyWFkBc9+zyBaKIqRrg/KwMqBbodBjgbHjDz7zjA=
github.com/beevik/etree v1.2.0 h1:l7WETslUG/T+xOPs47dtd6jov2Ii/8/OjCldk5fYfQw=
github.com/beevik/etree v1.2.0/go.mod h1:aiPf89g/1k3AShMVAzriilpcE4R/Vuor90y83zVZWFc=
github.com/benbjohnson/clock v1.3.0 h1:ip6w0uFQkncKQ979AypyG0ER7mqUSBdKLOgAle/AT8A=
github.com/benbjohnson/clock v1.3.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=
github.com/cheekybits/is v0.0.0-20150225183255-68e9c0620927 h1:SKI1/fuSdodxmNNyVBR8d7X/HuLnRpvvFO0AgyQk764=
@@ -17,8 +15,6 @@ github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs
github.com/decred/dcrd/crypto/blake256 v1.0.0 h1:/8DMNYp9SGi5f0w7uCm6d6M4OU2rGFK09Y2A4Xv7EE0=
github.com/decred/dcrd/dcrec/secp256k1/v4 v4.1.0 h1:HbphB4TFFXpv7MNrT52FGrrgVXF1owhMVTHFZIlnvd4=
github.com/decred/dcrd/dcrec/secp256k1/v4 v4.1.0/go.mod h1:DZGJHZMqrU4JJqFAWUS2UO1+lbSKsdiOoYi9Zzey7Fc=
github.com/go-ping/ping v1.1.0 h1:3MCGhVX4fyEUuhsfwPrsEdQw6xspHkv5zHsiSoDFZYw=
github.com/go-ping/ping v1.1.0/go.mod h1:xIFjORFzTxqIV/tDVGO4eDy/bLuSyawEeojSm3GfRGk=
github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA=
github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q=
github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q=
@@ -28,7 +24,6 @@ github.com/golang/protobuf v1.5.3/go.mod h1:XVQd3VNwM+JqD3oG2Ue2ip4fOMUkwXdXDdiu
github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE=
github.com/google/go-cmp v0.5.9 h1:O2Tfq5qg4qc4AmwVlvv0oLiVAGB7enBSJ2x2DqQFi38=
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/uuid v1.2.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/google/uuid v1.3.0 h1:t6JiXgmwXMjEs8VusXIJk2BXHsn+wx8BZdTaoZ5fu7I=
github.com/google/uuid v1.3.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gopherjs/gopherjs v1.17.2 h1:fQnZVsXk8uxXIStYb0N4bGk7jeyTalG/wsZjQ25dO0g=
@@ -147,25 +142,18 @@ golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU=
golang.org/x/net v0.0.0-20210316092652-d523dce5a7f4/go.mod h1:RBQZq4jEuRlivfhVLdyRGr576XBO4/greRjx4P4O3yc=
golang.org/x/net v0.8.0 h1:Zrh2ngAOFYneWTAIAPethzeaQLuHwhuBkuV6ZiRnUaQ=
golang.org/x/net v0.8.0/go.mod h1:QVkue5JL9kW//ek3r6jTKnTFis1tRmNAW2P1shuFdJc=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o=
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210315160823-c6e025ad8005/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20220704084225-05e143d24a9e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.6.0 h1:MVltZSvRTcU2ljQOhs94SXPftV6DCNnZViHeQps87pQ=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.8.0 h1:57P1ETyNKtuIjB4SRd15iJxuhj8Gc416Y78H3qgMh68=


+ 9
- 3
magefiles/common.go View File

@@ -12,8 +12,9 @@ import (
)

var Global = struct {
OS string
Arch string
OS string
Arch string
BuildRoot string
}{
Arch: "amd64",
}
@@ -30,7 +31,12 @@ type goBuildArgs struct {
}

func Build(args BuildArgs) error {
fullOutputDir, err := filepath.Abs(args.OutputDir)
buildRoot := Global.BuildRoot
if buildRoot == "" {
buildRoot = "build"
}

fullOutputDir, err := filepath.Abs(filepath.Join(buildRoot, args.OutputDir))
if err != nil {
return err
}


+ 5
- 0
magefiles/targets/targets.go View File

@@ -18,3 +18,8 @@ func Linux() {
func AMD64() {
magefiles.Global.Arch = "amd64"
}

// [配置项]设置编译的根目录
func BuildRoot(dir string) {
magefiles.Global.BuildRoot = dir
}

+ 0
- 37
models/models.go View File

@@ -24,40 +24,3 @@ func NewRepRedundancyConfig(repCount int) RepRedundancyConfig {

type ECRedundancyConfig struct {
}

// 冗余模式的具体配置
type RedundancyDataTypes interface{}
type RedundancyDataTypesConst interface {
RepRedundancyData | ECRedundancyData
}
type RepRedundancyData struct {
FileHash string `json:"fileHash"`
}

func NewRedundancyRepData(fileHash string) RepRedundancyData {
return RepRedundancyData{
FileHash: fileHash,
}
}

type ECRedundancyData struct {
Blocks []ObjectBlock `json:"blocks"`
}

func NewECRedundancyData(blocks []ObjectBlock) ECRedundancyData {
return ECRedundancyData{
Blocks: blocks,
}
}

type ObjectBlock struct {
Index int `json:"index"`
FileHash string `json:"fileHash"`
}

func NewObjectBlock(index int, fileHash string) ObjectBlock {
return ObjectBlock{
Index: index,
FileHash: fileHash,
}
}

pkg/actor/actor.go → pkgs/actor/actor.go View File

@@ -6,7 +6,7 @@ import (

"github.com/zyedidia/generic/list"

"gitlink.org.cn/cloudream/common/pkg/future"
"gitlink.org.cn/cloudream/common/pkgs/future"
mysync "gitlink.org.cn/cloudream/common/utils/sync"
)


pkg/actor/actor_test.go → pkgs/actor/actor_test.go View File


pkg/cmdtrie/command_trie.go → pkgs/cmdtrie/command_trie.go View File

@@ -6,7 +6,7 @@ import (
"strconv"

"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/pkg/trie"
"gitlink.org.cn/cloudream/common/pkgs/trie"
myreflect "gitlink.org.cn/cloudream/common/utils/reflect"
)


pkg/cmdtrie/command_trie_test.go → pkgs/cmdtrie/command_trie_test.go View File


pkg/distlock/config.go → pkgs/distlock/config.go View File


pkg/distlock/const.go → pkgs/distlock/const.go View File


pkg/distlock/distlock.go → pkgs/distlock/distlock.go View File


pkg/distlock/lockprovider/ipfs_lock.go → pkgs/distlock/lockprovider/ipfs_lock.go View File

@@ -4,7 +4,7 @@ import (
"fmt"

"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
mylo "gitlink.org.cn/cloudream/common/utils/lo"
)


pkg/distlock/lockprovider/ipfs_lock_test.go → pkgs/distlock/lockprovider/ipfs_lock_test.go View File

@@ -4,7 +4,7 @@ import (
"testing"

. "github.com/smartystreets/goconvey/convey"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
)

func Test_IPFSLock(t *testing.T) {

pkg/distlock/lockprovider/lock_compatibility_table.go → pkgs/distlock/lockprovider/lock_compatibility_table.go View File

@@ -4,7 +4,7 @@ import (
"fmt"

"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
)

const (

pkg/distlock/lockprovider/lock_compatibility_table_test.go → pkgs/distlock/lockprovider/lock_compatibility_table_test.go View File

@@ -4,7 +4,7 @@ import (
"testing"

. "github.com/smartystreets/goconvey/convey"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
)

func Test_LockCompatibilityTable(t *testing.T) {

pkg/distlock/lockprovider/metadata_lock.go → pkgs/distlock/lockprovider/metadata_lock.go View File

@@ -4,7 +4,7 @@ import (
"fmt"

"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
mylo "gitlink.org.cn/cloudream/common/utils/lo"
)


pkg/distlock/lockprovider/storage_lock.go → pkgs/distlock/lockprovider/storage_lock.go View File

@@ -4,7 +4,7 @@ import (
"fmt"

"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
mylo "gitlink.org.cn/cloudream/common/utils/lo"
)


pkg/distlock/lockprovider/string_lock_target.go → pkgs/distlock/lockprovider/string_lock_target.go View File


pkg/distlock/lockprovider/string_lock_target_test.go → pkgs/distlock/lockprovider/string_lock_target_test.go View File


pkg/distlock/reqbuilder/ipfs.go → pkgs/distlock/reqbuilder/ipfs.go View File

@@ -3,8 +3,8 @@ package reqbuilder
import (
"strconv"

"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type IPFSLockReqBuilder struct {

pkg/distlock/reqbuilder/lock_request_builder.go → pkgs/distlock/reqbuilder/lock_request_builder.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/service"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/service"
mylo "gitlink.org.cn/cloudream/common/utils/lo"
)


pkg/distlock/reqbuilder/metadata.go → pkgs/distlock/reqbuilder/metadata.go View File

@@ -1,6 +1,6 @@
package reqbuilder

import "gitlink.org.cn/cloudream/common/pkg/distlock"
import "gitlink.org.cn/cloudream/common/pkgs/distlock"

type MetadataLockReqBuilder struct {
*LockRequestBuilder

pkg/distlock/reqbuilder/metadata_bucket.go → pkgs/distlock/reqbuilder/metadata_bucket.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataBucketLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_cache.go → pkgs/distlock/reqbuilder/metadata_cache.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataCacheLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_node.go → pkgs/distlock/reqbuilder/metadata_node.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataNodeLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_object.go → pkgs/distlock/reqbuilder/metadata_object.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataObjectLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_object_block.go → pkgs/distlock/reqbuilder/metadata_object_block.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataObjectBlockLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_object_rep.go → pkgs/distlock/reqbuilder/metadata_object_rep.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataObjectRepLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_storage_object.go → pkgs/distlock/reqbuilder/metadata_storage_object.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataStorageObjectLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_user_bucket.go → pkgs/distlock/reqbuilder/metadata_user_bucket.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataUserBucketLockReqBuilder struct {

pkg/distlock/reqbuilder/metadata_user_storage.go → pkgs/distlock/reqbuilder/metadata_user_storage.go View File

@@ -1,8 +1,8 @@
package reqbuilder

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type MetadataUserStorageLockReqBuilder struct {

pkg/distlock/reqbuilder/storage.go → pkgs/distlock/reqbuilder/storage.go View File

@@ -3,8 +3,8 @@ package reqbuilder
import (
"strconv"

"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
)

type StorageLockReqBuilder struct {

pkg/distlock/service/init_providers.go → pkgs/distlock/service/init_providers.go View File

@@ -1,10 +1,10 @@
package service

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkg/distlock/service/internal"
"gitlink.org.cn/cloudream/common/pkg/trie"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/lockprovider"
"gitlink.org.cn/cloudream/common/pkgs/distlock/service/internal"
"gitlink.org.cn/cloudream/common/pkgs/trie"
)

func initProviders(providers *internal.ProvidersActor) {

pkg/distlock/service/internal/lease_actor.go → pkgs/distlock/service/internal/lease_actor.go View File

@@ -4,8 +4,8 @@ import (
"fmt"
"time"

"gitlink.org.cn/cloudream/common/pkg/actor"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkgs/actor"
"gitlink.org.cn/cloudream/common/pkgs/logger"
)

type lockRequestLease struct {

pkg/distlock/service/internal/main_actor.go → pkgs/distlock/service/internal/main_actor.go View File

@@ -6,8 +6,8 @@ import (
"strconv"
"time"

"gitlink.org.cn/cloudream/common/pkg/actor"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/actor"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/utils/serder"
clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"

pkg/distlock/service/internal/providers_actor.go → pkgs/distlock/service/internal/providers_actor.go View File

@@ -4,10 +4,10 @@ import (
"fmt"
"time"

"gitlink.org.cn/cloudream/common/pkg/actor"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/future"
"gitlink.org.cn/cloudream/common/pkg/trie"
"gitlink.org.cn/cloudream/common/pkgs/actor"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/future"
"gitlink.org.cn/cloudream/common/pkgs/trie"
)

type indexWaiter struct {

pkg/distlock/service/internal/retry_actor.go → pkgs/distlock/service/internal/retry_actor.go View File

@@ -5,10 +5,10 @@ import (
"time"

"github.com/samber/lo"
"gitlink.org.cn/cloudream/common/pkg/actor"
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/future"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkgs/actor"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/future"
"gitlink.org.cn/cloudream/common/pkgs/logger"
mylo "gitlink.org.cn/cloudream/common/utils/lo"
)


pkg/distlock/service/internal/utils.go → pkgs/distlock/service/internal/utils.go View File


pkg/distlock/service/internal/utils_test.go → pkgs/distlock/service/internal/utils_test.go View File


pkg/distlock/service/internal/watch_etcd_actor.go → pkgs/distlock/service/internal/watch_etcd_actor.go View File

@@ -4,7 +4,7 @@ import (
"context"
"fmt"

"gitlink.org.cn/cloudream/common/pkg/actor"
"gitlink.org.cn/cloudream/common/pkgs/actor"
mylo "gitlink.org.cn/cloudream/common/utils/lo"
"gitlink.org.cn/cloudream/common/utils/serder"
clientv3 "go.etcd.io/etcd/client/v3"

pkg/distlock/service/mutex.go → pkgs/distlock/service/mutex.go View File

@@ -1,7 +1,7 @@
package service

import (
"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
)

type Mutex struct {

pkg/distlock/service/service.go → pkgs/distlock/service/service.go View File

@@ -4,9 +4,9 @@ import (
"fmt"
"time"

"gitlink.org.cn/cloudream/common/pkg/distlock"
"gitlink.org.cn/cloudream/common/pkg/distlock/service/internal"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkgs/distlock"
"gitlink.org.cn/cloudream/common/pkgs/distlock/service/internal"
"gitlink.org.cn/cloudream/common/pkgs/logger"
clientv3 "go.etcd.io/etcd/client/v3"
)


pkg/event/event.go → pkgs/event/event.go View File


pkg/event/executor.go → pkgs/event/executor.go View File


pkg/future/future.go → pkgs/future/future.go View File


pkg/future/future_test.go → pkgs/future/future_test.go View File


pkg/future/set_value_future.go → pkgs/future/set_value_future.go View File


pkg/future/set_void_future.go → pkgs/future/set_void_future.go View File


pkg/logger/config.go → pkgs/logger/config.go View File


pkg/logger/global_logger.go → pkgs/logger/global_logger.go View File


pkg/logger/logger.go → pkgs/logger/logger.go View File


pkg/logger/logger_test.go → pkgs/logger/logger_test.go View File


pkg/logger/logrus_logger.go → pkgs/logger/logrus_logger.go View File


pkg/logger/utils.go → pkgs/logger/utils.go View File


pkg/mq/client.go → pkgs/mq/client.go View File

@@ -8,8 +8,8 @@ import (
"github.com/hashicorp/go-multierror"
"github.com/streadway/amqp"
"gitlink.org.cn/cloudream/common/consts/errorcode"
"gitlink.org.cn/cloudream/common/pkg/future"
"gitlink.org.cn/cloudream/common/pkg/logger"
"gitlink.org.cn/cloudream/common/pkgs/future"
"gitlink.org.cn/cloudream/common/pkgs/logger"
myreflect "gitlink.org.cn/cloudream/common/utils/reflect"
)


pkg/mq/message.go → pkgs/mq/message.go View File


pkg/mq/message_dispatcher.go → pkgs/mq/message_dispatcher.go View File


pkg/mq/message_test.go → pkgs/mq/message_test.go View File


pkg/mq/response.go → pkgs/mq/response.go View File


pkg/mq/server.go → pkgs/mq/server.go View File


pkg/task/manager.go → pkgs/task/manager.go View File


pkg/task/task.go → pkgs/task/task.go View File


pkg/tickevent/executor.go → pkgs/tickevent/executor.go View File


pkg/tickevent/tick_event.go → pkgs/tickevent/tick_event.go View File


pkg/trie/trie.go → pkgs/trie/trie.go View File


pkg/trie/trie_test.go → pkgs/trie/trie_test.go View File


pkg/typedispatcher/type_dispatcher.go → pkgs/typedispatcher/type_dispatcher.go View File


+ 0
- 74
utils/config.go View File

@@ -1,74 +0,0 @@
package utils

import (
"fmt"
"regexp"
"strconv"

"github.com/beevik/etree"
)

type EcConfig struct {
ecid string `xml:"ecid"`
class string `xml:"class"`
n int `xml:"n"`
k int `xml:"k"`
w int `xml:"w"`
opt int `xml:"opt"`
}

func (r *EcConfig) GetK() int {
return r.k
}

func (r *EcConfig) GetN() int {
return r.n
}

func GetEcPolicy() *map[string]EcConfig {
doc := etree.NewDocument()
if err := doc.ReadFromFile("../conf/sysSetting.xml"); err != nil {
panic(err)
}
ecMap := make(map[string]EcConfig, 20)
root := doc.SelectElement("setting")
for _, attr := range root.SelectElements("attribute") {
if name := attr.SelectElement("name"); name.Text() == "ec.policy" {
for _, eci := range attr.SelectElements("value") {
tt := EcConfig{}
tt.ecid = eci.SelectElement("ecid").Text()
tt.class = eci.SelectElement("class").Text()
tt.n, _ = strconv.Atoi(eci.SelectElement("n").Text())
tt.k, _ = strconv.Atoi(eci.SelectElement("k").Text())
tt.w, _ = strconv.Atoi(eci.SelectElement("w").Text())
tt.opt, _ = strconv.Atoi(eci.SelectElement("opt").Text())
ecMap[tt.ecid] = tt
}
}
}
fmt.Println(ecMap)
return &ecMap
//
}

func GetAgentIps() []string {
doc := etree.NewDocument()
if err := doc.ReadFromFile("../conf/sysSetting.xml"); err != nil {
panic(err)
}
root := doc.SelectElement("setting")
var ips []string // 定义存储 IP 的字符串切片

for _, attr := range root.SelectElements("attribute") {
if name := attr.SelectElement("name"); name.Text() == "agents.addr" {
for _, ip := range attr.SelectElements("value") {
ipRegex := regexp.MustCompile(`\b\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}\b`)
match := ipRegex.FindString(ip.Text())
print(match)
ips = append(ips, match)
}
}
}

return ips
}

+ 1
- 1
utils/config/config.go View File

@@ -28,7 +28,7 @@ func DefaultLoad(modeulName string, defCfg interface{}) error {
}

// TODO 可以考虑根据环境变量读取不同的配置
configFilePath := filepath.Join(filepath.Dir(execPath), "..", "conf", fmt.Sprintf("%s.config.json", modeulName))
configFilePath := filepath.Join(filepath.Dir(execPath), "..", "confs", fmt.Sprintf("%s.config.json", modeulName))
return Load(configFilePath, defCfg)
}



+ 0
- 123
utils/grpc/file_transport.go View File

@@ -1,123 +0,0 @@
package grpc

// TODO 拆分到存储服务的common包里去
/*
import (
"context"
"fmt"
"io"

myio "gitlink.org.cn/cloudream/common/utils/io"
"gitlink.org.cn/cloudream/proto"
)

type fileReadCloser struct {
io.ReadCloser
stream proto.FileTransport_GetFileClient
cancelFn context.CancelFunc
readData []byte
}

func (s *fileReadCloser) Read(p []byte) (int, error) {

if s.readData == nil {
resp, err := s.stream.Recv()
if err != nil {
return 0, err
}

if resp.Type == proto.FileDataPacketType_Data {
s.readData = resp.Data

} else if resp.Type == proto.FileDataPacketType_EOF {
return 0, io.EOF

} else {
return 0, fmt.Errorf("unsuppoted packt type: %v", resp.Type)
}
}

cnt := copy(p, s.readData)

if len(s.readData) == cnt {
s.readData = nil
} else {
s.readData = s.readData[cnt:]
}

return cnt, nil
}

func (s *fileReadCloser) Close() error {
s.cancelFn()

return nil
}

func GetFileAsStream(client proto.FileTransportClient, fileHash string) (io.ReadCloser, error) {
ctx, cancel := context.WithCancel(context.Background())

stream, err := client.GetFile(ctx, &proto.GetReq{
FileHash: fileHash,
})
if err != nil {
cancel()
return nil, fmt.Errorf("request grpc failed, err: %w", err)
}

return &fileReadCloser{
stream: stream,
cancelFn: cancel,
}, nil
}

type fileWriteCloser struct {
myio.PromiseWriteCloser[string]
stream proto.FileTransport_SendFileClient
}

func (s *fileWriteCloser) Write(p []byte) (int, error) {
err := s.stream.Send(&proto.FileDataPacket{
Type: proto.FileDataPacketType_Data,
Data: p,
})

if err != nil {
return 0, err
}

return len(p), nil
}

func (s *fileWriteCloser) Abort(err error) {
s.stream.CloseSend()
}

func (s *fileWriteCloser) Finish() (string, error) {
err := s.stream.Send(&proto.FileDataPacket{
Type: proto.FileDataPacketType_EOF,
})

if err != nil {
return "", fmt.Errorf("send EOF packet failed, err: %w", err)
}

resp, err := s.stream.CloseAndRecv()
if err != nil {
return "", fmt.Errorf("receive response failed, err: %w", err)
}

return resp.FileHash, nil
}

func SendFileAsStream(client proto.FileTransportClient) (myio.PromiseWriteCloser[string], error) {
stream, err := client.SendFile(context.Background())
if err != nil {
return nil, err
}

return &fileWriteCloser{
stream: stream,
}, nil
}
*/

+ 0
- 82
utils/ping.go View File

@@ -1,82 +0,0 @@
package utils

import (
//"fmt"
"github.com/go-ping/ping"
//"net"
"io/ioutil"
"net/http"
"strings"
"time"
)

type ConnStatus struct {
Addr string
IsReachable bool
Delay time.Duration
TTL int
}

// 获取本地主机 IP 地址
func getLocalIP() string {
resp, err := http.Get("https://api.ipify.org")
if err != nil {
panic(err)
}
defer resp.Body.Close()

body, err := ioutil.ReadAll(resp.Body)
if err != nil {
panic(err)
}

ip := strings.TrimSpace(string(body))
return ip
}

func GetConnStatus(remoteIP string) (*ConnStatus, error) {
// 本地主机 IP 地址
//localIP := getLocalIP()
//print("!@#@#!")
//print(localIP)
conn := ConnStatus{
Addr: remoteIP,
IsReachable: false,
}
pinger, err := ping.NewPinger(remoteIP)

if err != nil {
return nil, err
}
pinger.Count = 5 // 设置 ping 次数为 5
// pinger.Interval = 1 // 设置 ping 时间间隔为 1 秒
//pinger.Timeout = 2 // 设置 ping 超时时间为 2 秒
//pinger.SetPrivileged(true) // 设置使用特权模式以获取 TTL 值
pinger.OnRecv = func(pkt *ping.Packet) {
//fmt.Printf("%d bytes from %s: icmp_seq=%d time=%v ttl=%v (DUP!)\n",
// pkt.Nbytes, pkt.IPAddr, pkt.Seq, pkt.Rtt, pkt.Ttl)
conn.TTL = pkt.Ttl
}

/*pinger.OnDuplicateRecv = func(pkt *ping.Packet) {
fmt.Printf("%d bytes from %s: icmp_seq=%d time=%v ttl=%v (DUP!)\n",
pkt.Nbytes, pkt.IPAddr, pkt.Seq, pkt.Rtt, pkt.Ttl)
}*/

pinger.OnFinish = func(stats *ping.Statistics) {
//fmt.Printf("\n--- %s ping statistics ---\n", stats.Addr)
//fmt.Printf("%d packets transmitted, %d packets received, %v%% packet loss\n",
// stats.PacketsSent, stats.PacketsRecv, stats.PacketLoss)
//fmt.Printf("round-trip min/avg/max/stddev = %v/%v/%v/%v\n",
// stats.MinRtt, stats.AvgRtt, stats.MaxRtt, stats.StdDevRtt)
if stats.PacketLoss == 0.0 {
conn.IsReachable = true
}
conn.Delay = stats.AvgRtt
}
err = pinger.Run() // Blocks until finished.
if err != nil {
return nil, err
}
return &conn, nil
}

+ 0
- 21
utils/utils.go View File

@@ -1,21 +0,0 @@
package utils

import (
"fmt"
"strings"
)

// MakeMoveOperationFileName Move操作时,写入的文件的名称
func MakeMoveOperationFileName(objectID int64, userID int64) string {
return fmt.Sprintf("%d-%d", objectID, userID)
}

// GetDirectoryName 根据objectName获取所属的文件夹名
func GetDirectoryName(objectName string) string {
parts := strings.Split(objectName, "/")
//若为文件,dirName设置为空
if len(parts) == 1 {
return ""
}
return parts[0]
}

Loading…
Cancel
Save