| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859 |
- package mqttbroker
- import (
- "time"
- mqtt "github.com/mochi-mqtt/server/v2"
- "github.com/mochi-mqtt/server/v2/hooks/auth"
- "github.com/mochi-mqtt/server/v2/listeners"
- )
- // Broker 内嵌 MQTT broker,随应用进程启动,无需外部 mosquitto。
- // 局域网设备(相机/道闸)直连本 broker 的监听端口即可。
- type Broker struct {
- server *mqtt.Server
- serveErr chan error
- }
- // Start 启动内嵌 broker 并监听 addr(形如 ":1883" 或 "127.0.0.1:1883")。
- // 短暂等待 Serve 以捕获端口被占用等立即失败。
- func Start(addr string) (*Broker, error) {
- if addr == "" {
- addr = ":1883"
- }
- server := mqtt.New(&mqtt.Options{})
- // 局域网设备接入,允许匿名连接;如需鉴权再替换为 auth 鉴权 hook。
- if err := server.AddHook(new(auth.AllowHook), nil); err != nil {
- return nil, err
- }
- tcp := listeners.NewTCP(listeners.Config{
- ID: "t1",
- Address: addr,
- })
- if err := server.AddListener(tcp); err != nil {
- return nil, err
- }
- b := &Broker{server: server, serveErr: make(chan error, 1)}
- go func() {
- b.serveErr <- server.Serve()
- }()
- // 短暂等待,捕获 Serve 的立即失败(如端口被占用);正常运行不会返回。
- select {
- case err := <-b.serveErr:
- return nil, err
- case <-time.After(300 * time.Millisecond):
- return b, nil
- }
- }
- // Close 停止内嵌 broker。
- func (b *Broker) Close() error {
- if b == nil || b.server == nil {
- return nil
- }
- return b.server.Close()
- }
|