Нема описа
You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136
  1. package main
  2. import (
  3. "fmt"
  4. "log"
  5. "net/http"
  6. "os"
  7. "os/signal"
  8. "syscall"
  9. "git.x2erp.com/qdy/go-base/config"
  10. "git.x2erp.com/qdy/go-base/middleware"
  11. "git.x2erp.com/qdy/go-base/myservice"
  12. "git.x2erp.com/qdy/go-db/factory/rabbitmq"
  13. "git.x2erp.com/qdy/go-db/myhandle"
  14. "git.x2erp.com/qdy/go-svc-mqconsumer/functions"
  15. "go-micro.dev/v4/web"
  16. )
  17. func main() {
  18. cfg, err := config.GetConfig()
  19. if err != nil {
  20. log.Fatalf("Failed to create RabbitMQ factory: %v", err)
  21. return
  22. }
  23. serviceConfig := cfg.GetService()
  24. log.Printf("RabbitMQ Service Starting...")
  25. log.Printf("Service Port: %d", serviceConfig.Port)
  26. log.Printf("Service Name: %s", serviceConfig.ServiceName)
  27. // 启动微服务
  28. startRabbitMQService(cfg)
  29. }
  30. // 启动RabbitMQ微服务
  31. func startRabbitMQService(cfg config.IConfig) {
  32. // 初始化RabbitMQ工厂
  33. rabbitFactory, err := rabbitmq.NewRabbitMQFactory()
  34. if err != nil {
  35. log.Fatalf("Failed to create RabbitMQ factory: %v", err)
  36. }
  37. defer func() {
  38. if err := rabbitFactory.Close(); err != nil {
  39. log.Printf("RabbitMQ close error: %v", err)
  40. }
  41. }()
  42. // 设置优雅关闭
  43. setupGracefulShutdown(rabbitFactory)
  44. // 创建默认通道
  45. if err := createDefaultChannel(rabbitFactory); err != nil {
  46. log.Fatalf("Failed to create default channel: %v", err)
  47. }
  48. webService := myservice.StartStandalone(cfg)
  49. // 注册HTTP路由到webService
  50. registerRabbitMQRoutes(webService, rabbitFactory)
  51. // 等待服务运行
  52. log.Printf("RabbitMQ Service started successfully")
  53. log.Printf(" • Host: %s", cfg.GetRabbitMQ().Host)
  54. log.Printf(" • Port: %d", cfg.GetRabbitMQ().Port)
  55. log.Printf(" • Vhost: %s", cfg.GetRabbitMQ().Vhost)
  56. // 保持主程序运行
  57. select {}
  58. }
  59. // 创建默认通道
  60. func createDefaultChannel(rabbitFactory *rabbitmq.RabbitMQFactory) error {
  61. _, err := rabbitFactory.CreateChannel("default")
  62. if err != nil {
  63. return fmt.Errorf("failed to create default channel: %v", err)
  64. }
  65. log.Println("Default RabbitMQ channel created successfully")
  66. return nil
  67. }
  68. // 注册RabbitMQ相关路由
  69. func registerRabbitMQRoutes(webService web.Service, rabbitFactory *rabbitmq.RabbitMQFactory) {
  70. // 创建交换机
  71. webService.Handle("/api/rabbitmq/consumers/stop", middleware.JWTAuthMiddleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  72. myhandle.QueryHandlerJson(w, r, rabbitFactory, functions.CreateExchange)
  73. })))
  74. // 创建队列
  75. webService.Handle("/api/rabbitmq/consumers/start", middleware.JWTAuthMiddleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  76. myhandle.QueryHandlerJson(w, r, rabbitFactory, functions.CreateQueue)
  77. })))
  78. // 绑定队列到交换机
  79. webService.Handle("/api/rabbitmq/queue/bind", middleware.JWTAuthMiddleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  80. myhandle.QueryHandlerJson(w, r, rabbitFactory, functions.BindQueue)
  81. })))
  82. // 发送JSON消息
  83. webService.Handle("/api/rabbitmq/message/send", middleware.JWTAuthMiddleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  84. myhandle.QueryHandlerJson(w, r, rabbitFactory, functions.SendMessage)
  85. })))
  86. // 发送原始消息(字节流)
  87. webService.Handle("/api/rabbitmq/message/send/bytes", middleware.JWTAuthMiddleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  88. myhandle.QueryHandlerJson(w, r, rabbitFactory, functions.SendBytesMessage)
  89. })))
  90. // 健康检查
  91. webService.Handle("/api/rabbitmq/health", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  92. functions.HealthCheck(w, r, rabbitFactory)
  93. }))
  94. // 获取队列信息
  95. webService.Handle("/api/rabbitmq/queue/info", middleware.JWTAuthMiddleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  96. myhandle.QueryHandlerJson(w, r, rabbitFactory, functions.GetQueueInfo)
  97. })))
  98. }
  99. // 设置优雅关闭
  100. func setupGracefulShutdown(rabbitFactory *rabbitmq.RabbitMQFactory) {
  101. signalCh := make(chan os.Signal, 1)
  102. signal.Notify(signalCh, os.Interrupt, syscall.SIGTERM)
  103. go func() {
  104. <-signalCh
  105. log.Println("\nReceived shutdown signal, closing RabbitMQ connections...")
  106. if err := rabbitFactory.Close(); err != nil {
  107. log.Printf("Error closing RabbitMQ: %v", err)
  108. }
  109. log.Println("RabbitMQ connections closed gracefully")
  110. os.Exit(0)
  111. }()
  112. }