fix(bridge): resubscribe on every (re)connect; emqx v0.2.1 dashboard off; iot-service v0.1.1
This commit is contained in:
@@ -30,22 +30,18 @@ func Run(ctx context.Context, cfg *config.Config, sqsClient *sqs.Client, log *sl
|
||||
}
|
||||
log.Info("bridge: SQS queue resolved", "queue", cfg.SQSQueueName)
|
||||
|
||||
client, err := connectMQTT(ctx, cfg, log)
|
||||
// Обработчик сообщений передаём в connectMQTT: подписка выполняется в
|
||||
// OnConnectHandler, т.е. повторяется при каждом (ре)подключении.
|
||||
// Иначе после потери сессии EMQX бридж оставался бы без подписки
|
||||
// (прецедент 2026-08-16: реконнект без resubscribe → телеметрия терялась).
|
||||
handler := newMessageHandler(ctx, sqsClient, queueURL, log)
|
||||
|
||||
client, err := connectMQTT(ctx, cfg, log, handler)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer client.Disconnect(250)
|
||||
|
||||
handler := newMessageHandler(ctx, sqsClient, queueURL, log)
|
||||
token := client.Subscribe(telemetryTopicFilter, 1, handler)
|
||||
if !token.WaitTimeout(10 * time.Second) {
|
||||
return errSubscribeTimeout
|
||||
}
|
||||
if token.Error() != nil {
|
||||
return token.Error()
|
||||
}
|
||||
log.Info("bridge: subscribed", "filter", telemetryTopicFilter)
|
||||
|
||||
<-ctx.Done()
|
||||
log.Info("bridge: shutting down")
|
||||
return nil
|
||||
@@ -60,7 +56,8 @@ type errBridge string
|
||||
func (e errBridge) Error() string { return string(e) }
|
||||
|
||||
// connectMQTT устанавливает подключение к EMQX с автореконнектом.
|
||||
func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger) (mqtt.Client, error) {
|
||||
// handler передаётся в OnConnectHandler: подписка повторяется при каждом подключении.
|
||||
func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger, handler mqtt.MessageHandler) (mqtt.Client, error) {
|
||||
opts := mqtt.NewClientOptions()
|
||||
opts.AddBroker(cfg.MQTTBrokerURL)
|
||||
opts.SetClientID(cfg.MQTTClientID)
|
||||
@@ -78,8 +75,18 @@ func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger) (mqt
|
||||
opts.SetReconnectingHandler(func(_ mqtt.Client, _ *mqtt.ClientOptions) {
|
||||
log.Info("bridge: MQTT reconnecting...")
|
||||
})
|
||||
opts.SetOnConnectHandler(func(_ mqtt.Client) {
|
||||
opts.SetOnConnectHandler(func(c mqtt.Client) {
|
||||
log.Info("bridge: MQTT connected")
|
||||
token := c.Subscribe(telemetryTopicFilter, 1, handler)
|
||||
if !token.WaitTimeout(10 * time.Second) {
|
||||
log.Error("bridge: resubscribe timeout", "filter", telemetryTopicFilter)
|
||||
return
|
||||
}
|
||||
if token.Error() != nil {
|
||||
log.Error("bridge: resubscribe failed", "filter", telemetryTopicFilter, "err", token.Error())
|
||||
return
|
||||
}
|
||||
log.Info("bridge: subscribed", "filter", telemetryTopicFilter)
|
||||
})
|
||||
|
||||
client := mqtt.NewClient(opts)
|
||||
|
||||
Reference in New Issue
Block a user