server.go 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170
  1. package server
  2. import (
  3. "github.com/lf-edge/ekuiper/internal/conf"
  4. "github.com/lf-edge/ekuiper/internal/plugin"
  5. "github.com/lf-edge/ekuiper/internal/processor"
  6. "github.com/lf-edge/ekuiper/internal/service"
  7. "github.com/lf-edge/ekuiper/internal/xsql"
  8. "github.com/prometheus/client_golang/prometheus/promhttp"
  9. "context"
  10. "fmt"
  11. "net/http"
  12. "net/rpc"
  13. "os"
  14. "os/signal"
  15. "path"
  16. "syscall"
  17. "time"
  18. )
  19. var (
  20. dataDir string
  21. logger = conf.Log
  22. startTimeStamp int64
  23. version = ""
  24. ruleProcessor *processor.RuleProcessor
  25. streamProcessor *processor.StreamProcessor
  26. pluginManager *plugin.Manager
  27. serviceManager *service.Manager
  28. )
  29. func StartUp(Version, LoadFileType string) {
  30. version = Version
  31. conf.LoadFileType = LoadFileType
  32. startTimeStamp = time.Now().Unix()
  33. conf.InitConf()
  34. dr, err := conf.GetDataLoc()
  35. if err != nil {
  36. logger.Panic(err)
  37. } else {
  38. logger.Infof("db location is %s", dr)
  39. dataDir = dr
  40. }
  41. ruleProcessor = processor.NewRuleProcessor(dataDir)
  42. streamProcessor = processor.NewStreamProcessor(path.Join(dataDir, "stream"))
  43. pluginManager, err = plugin.NewPluginManager()
  44. if err != nil {
  45. logger.Panic(err)
  46. }
  47. serviceManager, err = service.GetServiceManager()
  48. if err != nil {
  49. logger.Panic(err)
  50. }
  51. xsql.InitFuncRegisters(serviceManager, pluginManager)
  52. registry = &RuleRegistry{internal: make(map[string]*RuleState)}
  53. server := new(Server)
  54. //Start rules
  55. if rules, err := ruleProcessor.GetAllRules(); err != nil {
  56. logger.Infof("Start rules error: %s", err)
  57. } else {
  58. logger.Info("Starting rules")
  59. var reply string
  60. for _, rule := range rules {
  61. //err = server.StartRule(rule, &reply)
  62. reply = recoverRule(rule)
  63. if 0 != len(reply) {
  64. logger.Info(reply)
  65. }
  66. }
  67. }
  68. //Start prometheus service
  69. var srvPrometheus *http.Server = nil
  70. if conf.Config.Basic.Prometheus {
  71. portPrometheus := conf.Config.Basic.PrometheusPort
  72. if portPrometheus <= 0 {
  73. logger.Fatal("Miss configuration prometheusPort")
  74. }
  75. mux := http.NewServeMux()
  76. mux.Handle("/metrics", promhttp.Handler())
  77. srvPrometheus = &http.Server{
  78. Addr: fmt.Sprintf("0.0.0.0:%d", portPrometheus),
  79. WriteTimeout: time.Second * 15,
  80. ReadTimeout: time.Second * 15,
  81. IdleTimeout: time.Second * 60,
  82. Handler: mux,
  83. }
  84. go func() {
  85. if err := srvPrometheus.ListenAndServe(); err != nil && err != http.ErrServerClosed {
  86. logger.Fatal("Listen prometheus error: ", err)
  87. }
  88. }()
  89. msg := fmt.Sprintf("Serving prometheus metrics on port http://localhost:%d/metrics", portPrometheus)
  90. logger.Infof(msg)
  91. fmt.Println(msg)
  92. }
  93. //Start rest service
  94. srvRest := createRestServer(conf.Config.Basic.RestIp, conf.Config.Basic.RestPort)
  95. go func() {
  96. var err error
  97. if conf.Config.Basic.RestTls == nil {
  98. err = srvRest.ListenAndServe()
  99. } else {
  100. err = srvRest.ListenAndServeTLS(conf.Config.Basic.RestTls.Certfile, conf.Config.Basic.RestTls.Keyfile)
  101. }
  102. if err != nil && err != http.ErrServerClosed {
  103. logger.Fatal("Error serving rest service: ", err)
  104. }
  105. }()
  106. // Start rpc service
  107. portRpc := conf.Config.Basic.Port
  108. ipRpc := conf.Config.Basic.Ip
  109. rpcSrv := rpc.NewServer()
  110. err = rpcSrv.Register(server)
  111. if err != nil {
  112. logger.Fatal("Format of service Server isn'restHttpType correct. ", err)
  113. }
  114. srvRpc := &http.Server{
  115. Addr: fmt.Sprintf("%s:%d", ipRpc, portRpc),
  116. WriteTimeout: time.Second * 15,
  117. ReadTimeout: time.Second * 15,
  118. IdleTimeout: time.Second * 60,
  119. Handler: rpcSrv,
  120. }
  121. go func() {
  122. if err = srvRpc.ListenAndServe(); err != nil && err != http.ErrServerClosed {
  123. logger.Fatal("Error serving rpc service:", err)
  124. }
  125. }()
  126. //Startup message
  127. restHttpType := "http"
  128. if conf.Config.Basic.RestTls != nil {
  129. restHttpType = "https"
  130. }
  131. msg := fmt.Sprintf("Serving kuiper (version - %s) on port %d, and restful api on %s://%s:%d. \n", Version, conf.Config.Basic.Port, restHttpType, conf.Config.Basic.RestIp, conf.Config.Basic.RestPort)
  132. logger.Info(msg)
  133. fmt.Printf(msg)
  134. //Stop the services
  135. sigint := make(chan os.Signal, 1)
  136. signal.Notify(sigint, os.Interrupt, syscall.SIGTERM)
  137. <-sigint
  138. if err = srvRpc.Shutdown(context.TODO()); err != nil {
  139. logger.Errorf("rpc server shutdown error: %v", err)
  140. }
  141. logger.Info("rpc server shutdown.")
  142. if err = srvRest.Shutdown(context.TODO()); err != nil {
  143. logger.Errorf("rest server shutdown error: %v", err)
  144. }
  145. logger.Info("rest server successfully shutdown.")
  146. if srvPrometheus != nil {
  147. if err = srvPrometheus.Shutdown(context.TODO()); err != nil {
  148. logger.Errorf("prometheus server shutdown error: %v", err)
  149. }
  150. logger.Info("prometheus server successfully shutdown.")
  151. }
  152. os.Exit(0)
  153. }