説明なし
選択できるのは25トピックまでです。 トピックは、先頭が英数字で、英数字とダッシュ('-')を使用した35文字以内のものにしてください。

main.go 5.6KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221
  1. package main
  2. import (
  3. "encoding/json"
  4. "log"
  5. "net/http"
  6. "strings"
  7. "time"
  8. "git.x2erp.com/qdy/go-base/config"
  9. "git.x2erp.com/qdy/go-base/logger"
  10. "git.x2erp.com/qdy/go-base/myservice"
  11. "git.x2erp.com/qdy/go-base/types"
  12. "git.x2erp.com/qdy/go-db/factory/database"
  13. "git.x2erp.com/qdy/go-svc-worker/service"
  14. "go-micro.dev/v4/metadata"
  15. )
  16. // 定义业务服务
  17. type DBFactory struct {
  18. dbFactory *database.DBFactory
  19. }
  20. // 配置
  21. var (
  22. cfg config.IConfig
  23. serviceName string
  24. serviceVersion string
  25. )
  26. func main() {
  27. // ========== 第一阶段:强制写入文件的启动日志 ==========
  28. // 这一步确保即使配置加载失败,也有日志记录
  29. if err := logger.InitBootLog("svc-worker"); err != nil {
  30. // 连启动日志都初始化失败,只能输出到控制台
  31. log.Fatal("无法初始化启动日志: ", err)
  32. }
  33. logger.BootLog("开始加载配置...")
  34. cfg, err := config.GetConfig()
  35. if err != nil {
  36. log.Fatalf("Failed to create RabbitMQ factory: %v", err)
  37. }
  38. serviceConfig := cfg.GetService()
  39. microConfig := cfg.GetMicro()
  40. serviceName = serviceConfig.ServiceName
  41. serviceVersion = serviceConfig.ServiceVersion
  42. log.Printf("serviceName: %s", serviceName)
  43. log.Printf("Port: %d", serviceConfig.Port)
  44. log.Printf("Consul: %s", microConfig.RegistryAddress)
  45. // 2. 初始化数据库
  46. dbFactory, err := database.GetDBFactory()
  47. if err != nil {
  48. log.Fatal("数据库连接失败:", err)
  49. }
  50. defer func() {
  51. if err := dbFactory.Close(); err != nil {
  52. logger.Info("数据库关闭错误: %v", err)
  53. }
  54. }()
  55. // 3. 创建服务实例
  56. dbfactory := &DBFactory{dbFactory: dbFactory}
  57. // 4. 使用 micro.Start 启动服务
  58. webService := myservice.Start(cfg)
  59. // 7. 注册HTTP路由
  60. webService.Handle("/", http.HandlerFunc(rootHandler))
  61. webService.Handle("/health", http.HandlerFunc(dbfactory.healthHandler))
  62. webService.Handle("/info", http.HandlerFunc(infoHandler))
  63. webService.Handle("/api/data/agent/to/doris", authMiddleware(http.HandlerFunc(dbfactory.agentToDorisHandler)))
  64. logger.InitRuntimeLogger("order-service", cfg.GetLog())
  65. log.Println("日志系统初始化完成")
  66. //关闭-启动日志输出文件功能
  67. logger.CloseBootLogger()
  68. if err := webService.Run(); err != nil {
  69. log.Fatal("服务运行失败:", err)
  70. }
  71. }
  72. // 根处理器
  73. func rootHandler(w http.ResponseWriter, r *http.Request) {
  74. if r.URL.Path != "/" {
  75. http.NotFound(w, r)
  76. return
  77. }
  78. respondJSON(w, http.StatusOK, map[string]string{
  79. "service": serviceName,
  80. "status": "running",
  81. "mode": "http-microservice",
  82. })
  83. }
  84. // 健康检查处理器
  85. func (s *DBFactory) healthHandler(w http.ResponseWriter, r *http.Request) {
  86. if err := s.dbFactory.TestConnection(s.dbFactory.GetDBType()); err != nil {
  87. respondJSON(w, http.StatusServiceUnavailable, map[string]string{
  88. "status": "down",
  89. "error": err.Error(),
  90. })
  91. return
  92. }
  93. respondJSON(w, http.StatusOK, map[string]string{
  94. "status": "up",
  95. "time": time.Now().Format(time.RFC3339),
  96. })
  97. }
  98. // 信息处理器
  99. func infoHandler(w http.ResponseWriter, r *http.Request) {
  100. respondJSON(w, http.StatusOK, map[string]interface{}{
  101. "service": serviceName,
  102. "version": serviceVersion,
  103. "api": map[string]string{
  104. "POST /api/data/agent/to/doris": "同步数据到Doris",
  105. "GET /health": "健康检查",
  106. "GET /info": "服务信息",
  107. "GET /": "根路径",
  108. },
  109. "features": []string{
  110. "服务发现(Consul)",
  111. "负载均衡",
  112. "健康检查",
  113. "HTTP API网关",
  114. },
  115. })
  116. }
  117. // AgentToDoris处理器
  118. func (s *DBFactory) agentToDorisHandler(w http.ResponseWriter, r *http.Request) {
  119. if r.Method != "POST" {
  120. respondJSON(w, http.StatusMethodNotAllowed, types.QueryResult{
  121. Error: "只支持POST请求",
  122. Success: false,
  123. })
  124. return
  125. }
  126. // 解析请求
  127. var requestData types.QueryRequest
  128. if err := json.NewDecoder(r.Body).Decode(&requestData); err != nil {
  129. respondJSON(w, http.StatusBadRequest, map[string]string{
  130. "error": "无效的JSON数据",
  131. })
  132. return
  133. }
  134. // 处理业务逻辑
  135. result := service.ServiceAgentToDoris(s.dbFactory, requestData)
  136. respondJSON(w, http.StatusOK, result)
  137. }
  138. // 认证中间件
  139. func authMiddleware(next http.Handler) http.Handler {
  140. return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  141. // JWT令牌认证
  142. token := r.Header.Get("Authorization")
  143. if token != "" && strings.HasPrefix(token, "Bearer ") {
  144. token = token[7:]
  145. }
  146. // 双重认证:API密钥或JWT
  147. if token == "" {
  148. respondJSON(w, http.StatusUnauthorized, map[string]string{
  149. "error": "需要API密钥或Bearer令牌",
  150. })
  151. return
  152. }
  153. // 验证JWT令牌
  154. if token != "" && !isValidJWT(token) {
  155. respondJSON(w, http.StatusUnauthorized, map[string]string{
  156. "error": "无效的访问令牌",
  157. })
  158. return
  159. }
  160. // 将认证信息添加到上下文
  161. ctx := r.Context()
  162. if token != "" {
  163. ctx = metadata.Set(ctx, "Authorization", "Bearer "+token)
  164. }
  165. next.ServeHTTP(w, r.WithContext(ctx))
  166. })
  167. }
  168. // JWT验证
  169. func isValidJWT(token string) bool {
  170. // TODO: 实现JWT验证逻辑
  171. // 可以使用 github.com/golang-jwt/jwt/v5
  172. // 临时实现:检查token是否有效格式
  173. //if len(token) < 10 {
  174. // return false
  175. // }
  176. return true // 临时返回true,实际需要验证签名和过期时间
  177. }
  178. // JSON响应辅助函数
  179. func respondJSON(w http.ResponseWriter, status int, data interface{}) {
  180. w.Header().Set("Content-Type", "application/json")
  181. w.Header().Set("X-Service-Name", serviceName)
  182. w.Header().Set("X-Service-Version", serviceVersion)
  183. w.WriteHeader(status)
  184. if err := json.NewEncoder(w).Encode(data); err != nil {
  185. log.Printf("JSON编码错误: %v", err)
  186. }
  187. }