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.

ops.go 2.4 kB

2 years ago
2 years ago
2 years ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134
  1. package ops
  2. import (
  3. "fmt"
  4. "gitlink.org.cn/cloudream/common/pkgs/ioswitch/dag"
  5. "gitlink.org.cn/cloudream/common/pkgs/ioswitch/exec"
  6. "gitlink.org.cn/cloudream/common/pkgs/types"
  7. cdssdk "gitlink.org.cn/cloudream/common/sdks/storage"
  8. "gitlink.org.cn/cloudream/common/utils/serder"
  9. )
  10. var OpUnion = serder.UseTypeUnionExternallyTagged(types.Ref(types.NewTypeUnion[exec.Op]()))
  11. type AgentWorker struct {
  12. Node cdssdk.Node
  13. }
  14. func (w *AgentWorker) GetAddress() string {
  15. // TODO 选择地址
  16. return fmt.Sprintf("%v:%v", w.Node.ExternalIP, w.Node.ExternalGRPCPort)
  17. }
  18. func (w *AgentWorker) Equals(worker dag.WorkerInfo) bool {
  19. aw, ok := worker.(*AgentWorker)
  20. if !ok {
  21. return false
  22. }
  23. return w.Node.NodeID == aw.Node.NodeID
  24. }
  25. type NodeProps struct {
  26. From From
  27. To To
  28. }
  29. type ValueVarType int
  30. const (
  31. StringValueVar ValueVarType = iota
  32. SignalValueVar
  33. )
  34. type VarProps struct {
  35. StreamIndex int // 流的编号,只在StreamVar上有意义
  36. ValueType ValueVarType // 值类型,只在ValueVar上有意义
  37. Var exec.Var // 生成Plan的时候创建的对应的Var
  38. }
  39. type Graph = dag.Graph[NodeProps, VarProps]
  40. type Node = dag.Node[NodeProps, VarProps]
  41. type StreamVar = dag.StreamVar[NodeProps, VarProps]
  42. type ValueVar = dag.ValueVar[NodeProps, VarProps]
  43. func addOpByEnv(op exec.Op, env dag.NodeEnv, blder *exec.PlanBuilder) {
  44. switch env.Type {
  45. case dag.EnvWorker:
  46. blder.AtAgent(env.Worker.(*AgentWorker).Node).AddOp(op)
  47. case dag.EnvExecutor:
  48. blder.AtExecutor().AddOp(op)
  49. }
  50. }
  51. func formatStreamIO(node *Node) string {
  52. is := ""
  53. for i, in := range node.InputStreams {
  54. if i > 0 {
  55. is += ","
  56. }
  57. if in == nil {
  58. is += "."
  59. } else {
  60. is += fmt.Sprintf("%v", in.ID)
  61. }
  62. }
  63. os := ""
  64. for i, out := range node.OutputStreams {
  65. if i > 0 {
  66. os += ","
  67. }
  68. if out == nil {
  69. os += "."
  70. } else {
  71. os += fmt.Sprintf("%v", out.ID)
  72. }
  73. }
  74. if is == "" && os == "" {
  75. return ""
  76. }
  77. return fmt.Sprintf("S{%s>%s}", is, os)
  78. }
  79. func formatValueIO(node *Node) string {
  80. is := ""
  81. for i, in := range node.InputValues {
  82. if i > 0 {
  83. is += ","
  84. }
  85. if in == nil {
  86. is += "."
  87. } else {
  88. is += fmt.Sprintf("%v", in.ID)
  89. }
  90. }
  91. os := ""
  92. for i, out := range node.OutputValues {
  93. if i > 0 {
  94. os += ","
  95. }
  96. if out == nil {
  97. os += "."
  98. } else {
  99. os += fmt.Sprintf("%v", out.ID)
  100. }
  101. }
  102. if is == "" && os == "" {
  103. return ""
  104. }
  105. return fmt.Sprintf("V{%s>%s}", is, os)
  106. }

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