broker.go 1.4 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859
  1. package mqttbroker
  2. import (
  3. "time"
  4. mqtt "github.com/mochi-mqtt/server/v2"
  5. "github.com/mochi-mqtt/server/v2/hooks/auth"
  6. "github.com/mochi-mqtt/server/v2/listeners"
  7. )
  8. // Broker 内嵌 MQTT broker,随应用进程启动,无需外部 mosquitto。
  9. // 局域网设备(相机/道闸)直连本 broker 的监听端口即可。
  10. type Broker struct {
  11. server *mqtt.Server
  12. serveErr chan error
  13. }
  14. // Start 启动内嵌 broker 并监听 addr(形如 ":1883" 或 "127.0.0.1:1883")。
  15. // 短暂等待 Serve 以捕获端口被占用等立即失败。
  16. func Start(addr string) (*Broker, error) {
  17. if addr == "" {
  18. addr = ":1883"
  19. }
  20. server := mqtt.New(&mqtt.Options{})
  21. // 局域网设备接入,允许匿名连接;如需鉴权再替换为 auth 鉴权 hook。
  22. if err := server.AddHook(new(auth.AllowHook), nil); err != nil {
  23. return nil, err
  24. }
  25. tcp := listeners.NewTCP(listeners.Config{
  26. ID: "t1",
  27. Address: addr,
  28. })
  29. if err := server.AddListener(tcp); err != nil {
  30. return nil, err
  31. }
  32. b := &Broker{server: server, serveErr: make(chan error, 1)}
  33. go func() {
  34. b.serveErr <- server.Serve()
  35. }()
  36. // 短暂等待,捕获 Serve 的立即失败(如端口被占用);正常运行不会返回。
  37. select {
  38. case err := <-b.serveErr:
  39. return nil, err
  40. case <-time.After(300 * time.Millisecond):
  41. return b, nil
  42. }
  43. }
  44. // Close 停止内嵌 broker。
  45. func (b *Broker) Close() error {
  46. if b == nil || b.server == nil {
  47. return nil
  48. }
  49. return b.server.Close()
  50. }