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.

main.go 3.7 kB

2 years ago
2 years ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130
  1. package main
  2. import (
  3. "fmt"
  4. "os"
  5. "gitlink.org.cn/cloudream/common/pkgs/logger"
  6. stgglb "gitlink.org.cn/cloudream/storage/common/globals"
  7. "gitlink.org.cn/cloudream/storage/common/pkgs/db"
  8. "gitlink.org.cn/cloudream/storage/common/pkgs/distlock"
  9. scmq "gitlink.org.cn/cloudream/storage/common/pkgs/mq/scanner"
  10. "gitlink.org.cn/cloudream/storage/scanner/internal/config"
  11. "gitlink.org.cn/cloudream/storage/scanner/internal/event"
  12. "gitlink.org.cn/cloudream/storage/scanner/internal/mq"
  13. "gitlink.org.cn/cloudream/storage/scanner/internal/tickevent"
  14. )
  15. func main() {
  16. err := config.Init()
  17. if err != nil {
  18. fmt.Printf("init config failed, err: %s", err.Error())
  19. os.Exit(1)
  20. }
  21. err = logger.Init(&config.Cfg().Logger)
  22. if err != nil {
  23. fmt.Printf("init logger failed, err: %s", err.Error())
  24. os.Exit(1)
  25. }
  26. db, err := db.NewDB(&config.Cfg().DB)
  27. if err != nil {
  28. logger.Fatalf("new db failed, err: %s", err.Error())
  29. }
  30. stgglb.InitMQPool(&config.Cfg().RabbitMQ)
  31. distlockSvc, err := distlock.NewService(&config.Cfg().DistLock)
  32. if err != nil {
  33. logger.Warnf("new distlock service failed, err: %s", err.Error())
  34. os.Exit(1)
  35. }
  36. go serveDistLock(distlockSvc)
  37. eventExecutor := event.NewExecutor(db, distlockSvc)
  38. go serveEventExecutor(&eventExecutor)
  39. agtSvr, err := scmq.NewServer(mq.NewService(&eventExecutor), &config.Cfg().RabbitMQ)
  40. if err != nil {
  41. logger.Fatalf("new agent server failed, err: %s", err.Error())
  42. }
  43. agtSvr.OnError(func(err error) {
  44. logger.Warnf("agent server err: %s", err.Error())
  45. })
  46. go serveScannerServer(agtSvr)
  47. tickExecutor := tickevent.NewExecutor(tickevent.ExecuteArgs{
  48. EventExecutor: &eventExecutor,
  49. DB: db,
  50. })
  51. startTickEvent(&tickExecutor)
  52. forever := make(chan struct{})
  53. <-forever
  54. }
  55. func serveEventExecutor(executor *event.Executor) {
  56. logger.Info("start serving event executor")
  57. err := executor.Execute()
  58. if err != nil {
  59. logger.Errorf("event executor stopped with error: %s", err.Error())
  60. }
  61. logger.Info("event executor stopped")
  62. // TODO 仅简单结束了程序
  63. os.Exit(1)
  64. }
  65. func serveScannerServer(server *scmq.Server) {
  66. logger.Info("start serving scanner server")
  67. err := server.Serve()
  68. if err != nil {
  69. logger.Errorf("scanner server stopped with error: %s", err.Error())
  70. }
  71. logger.Info("scanner server stopped")
  72. // TODO 仅简单结束了程序
  73. os.Exit(1)
  74. }
  75. func serveDistLock(svc *distlock.Service) {
  76. logger.Info("start serving distlock")
  77. err := svc.Serve()
  78. if err != nil {
  79. logger.Errorf("distlock stopped with error: %s", err.Error())
  80. }
  81. logger.Info("distlock stopped")
  82. // TODO 仅简单结束了程序
  83. os.Exit(1)
  84. }
  85. func startTickEvent(tickExecutor *tickevent.Executor) {
  86. // TODO 可以考虑增加配置文件,配置这些任务间隔时间
  87. interval := 5 * 60 * 1000
  88. tickExecutor.Start(tickevent.NewBatchAllAgentCheckCache(), interval, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  89. tickExecutor.Start(tickevent.NewBatchCheckAllPackage(), interval, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  90. // tickExecutor.Start(tickevent.NewBatchCheckAllRepCount(), interval, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  91. tickExecutor.Start(tickevent.NewBatchCheckAllStorage(), interval, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  92. tickExecutor.Start(tickevent.NewCheckAgentState(), 5*60*1000, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  93. tickExecutor.Start(tickevent.NewBatchCheckPackageRedundancy(), interval, tickevent.StartOption{RandomStartDelayMs: 20 * 60 * 1000})
  94. tickExecutor.Start(tickevent.NewBatchCleanPinned(), interval, tickevent.StartOption{RandomStartDelayMs: 20 * 60 * 1000})
  95. }

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