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.0 kB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108
  1. package main
  2. import (
  3. "fmt"
  4. "os"
  5. "sync"
  6. distlocksvc "gitlink.org.cn/cloudream/common/pkg/distlock/service"
  7. log "gitlink.org.cn/cloudream/common/pkg/logger"
  8. "gitlink.org.cn/cloudream/db"
  9. scsvr "gitlink.org.cn/cloudream/rabbitmq/server/scanner"
  10. "gitlink.org.cn/cloudream/scanner/internal/config"
  11. "gitlink.org.cn/cloudream/scanner/internal/event"
  12. "gitlink.org.cn/cloudream/scanner/internal/services"
  13. "gitlink.org.cn/cloudream/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 = log.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. log.Fatalf("new db failed, err: %s", err.Error())
  29. }
  30. distlockSvc, err := distlocksvc.NewService(&config.Cfg().DistLock)
  31. if err != nil {
  32. log.Warnf("new distlock service failed, err: %s", err.Error())
  33. os.Exit(1)
  34. }
  35. wg := sync.WaitGroup{}
  36. wg.Add(2)
  37. eventExecutor := event.NewExecutor(db, distlockSvc)
  38. go serveEventExecutor(&eventExecutor, &wg)
  39. agtSvr, err := scsvr.NewServer(services.NewService(&eventExecutor), &config.Cfg().RabbitMQ)
  40. if err != nil {
  41. log.Fatalf("new agent server failed, err: %s", err.Error())
  42. }
  43. agtSvr.OnError = func(err error) {
  44. log.Warnf("agent server err: %s", err.Error())
  45. }
  46. go serveScannerServer(agtSvr, &wg)
  47. tickExecutor := tickevent.NewExecutor(tickevent.ExecuteArgs{
  48. EventExecutor: &eventExecutor,
  49. DB: db,
  50. })
  51. startTickEvent(&tickExecutor)
  52. wg.Wait()
  53. }
  54. func serveEventExecutor(executor *event.Executor, wg *sync.WaitGroup) {
  55. log.Info("start serving event executor")
  56. err := executor.Execute()
  57. if err != nil {
  58. log.Errorf("event executor stopped with error: %s", err.Error())
  59. }
  60. log.Info("event executor stopped")
  61. wg.Done()
  62. }
  63. func serveScannerServer(server *scsvr.Server, wg *sync.WaitGroup) {
  64. log.Info("start serving scanner server")
  65. err := server.Serve()
  66. if err != nil {
  67. log.Errorf("scanner server stopped with error: %s", err.Error())
  68. }
  69. log.Info("scanner server stopped")
  70. wg.Done()
  71. }
  72. func startTickEvent(tickExecutor *tickevent.Executor) {
  73. // TODO 可以考虑增加配置文件,配置这些任务间隔时间
  74. tickExecutor.Start(tickevent.NewBatchAllAgentCheckCache(), 5*60*1000, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  75. tickExecutor.Start(tickevent.NewBatchCheckAllObject(), 5*60*1000, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  76. tickExecutor.Start(tickevent.NewBatchCheckAllRepCount(), 5*60*1000, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  77. tickExecutor.Start(tickevent.NewBatchCheckAllStorage(), 5*60*1000, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  78. //tickExecutor.Start(tickevent.NewCheckAgentState(), 5*60*1000, tickevent.StartOption{RandomFirstStartDelayMs: 60 * 1000})
  79. tickExecutor.Start(tickevent.NewCheckCache(), 5*60*1000, tickevent.StartOption{RandomStartDelayMs: 60 * 1000})
  80. }

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