safety fix
This commit is contained in:
parent
a872efdb2e
commit
39365d1b08
1 changed files with 32 additions and 5 deletions
|
|
@ -129,8 +129,11 @@ func (c *MQTTChannel) Start(ctx context.Context) error {
|
||||||
logger.InfoC("mqtt", "Connected to MQTT broker")
|
logger.InfoC("mqtt", "Connected to MQTT broker")
|
||||||
|
|
||||||
// Subscribe to topics after successful connection
|
// Subscribe to topics after successful connection
|
||||||
|
var subscriptionErrors []string
|
||||||
for _, topic := range c.config.SubscribeTopics {
|
for _, topic := range c.config.SubscribeTopics {
|
||||||
if token := c.client.Subscribe(topic, byte(c.config.QoS), nil); token.Wait() && token.Error() != nil {
|
if token := c.client.Subscribe(topic, byte(c.config.QoS), nil); token.Wait() && token.Error() != nil {
|
||||||
|
errMsg := fmt.Sprintf("topic %s: %v", topic, token.Error())
|
||||||
|
subscriptionErrors = append(subscriptionErrors, errMsg)
|
||||||
logger.ErrorCF("mqtt", "Failed to subscribe to topic", map[string]any{
|
logger.ErrorCF("mqtt", "Failed to subscribe to topic", map[string]any{
|
||||||
"topic": topic,
|
"topic": topic,
|
||||||
"error": token.Error(),
|
"error": token.Error(),
|
||||||
|
|
@ -141,12 +144,22 @@ func (c *MQTTChannel) Start(ctx context.Context) error {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Log subscription errors summary
|
||||||
|
if len(subscriptionErrors) > 0 {
|
||||||
|
logger.ErrorCF("mqtt", "Subscription errors occurred", map[string]any{
|
||||||
|
"errors": subscriptionErrors,
|
||||||
|
"count": len(subscriptionErrors),
|
||||||
|
})
|
||||||
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
c.client = mqtt.NewClient(opts)
|
c.client = mqtt.NewClient(opts)
|
||||||
|
|
||||||
// Connect with retry
|
// Connect with retry
|
||||||
if token := c.client.Connect(); token.Wait() && token.Error() != nil {
|
if token := c.client.Connect(); token.Wait() && token.Error() != nil {
|
||||||
|
// Clean up client on connection failure
|
||||||
|
c.client = nil
|
||||||
return fmt.Errorf("mqtt connect failed: %w", token.Error())
|
return fmt.Errorf("mqtt connect failed: %w", token.Error())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -236,11 +249,18 @@ func (c *MQTTChannel) onMessage(client mqtt.Client, msg mqtt.Message) {
|
||||||
// Check if message has reply-to header (in payload or topic)
|
// Check if message has reply-to header (in payload or topic)
|
||||||
replyTopic := c.config.ReplyTopic
|
replyTopic := c.config.ReplyTopic
|
||||||
if strings.Contains(content, "reply-to:") {
|
if strings.Contains(content, "reply-to:") {
|
||||||
// Simple parsing for reply-to
|
// Improved parsing for reply-to: extract the last occurrence
|
||||||
parts := strings.SplitN(content, "reply-to:", 2)
|
lastIndex := strings.LastIndex(content, "reply-to:")
|
||||||
if len(parts) == 2 {
|
if lastIndex != -1 {
|
||||||
replyTopic = strings.TrimSpace(parts[1])
|
replyTopicPart := content[lastIndex+len("reply-to:"):]
|
||||||
content = strings.TrimSpace(parts[0])
|
// Find the end of the reply-to value (newline or end of string)
|
||||||
|
if newlineIndex := strings.Index(replyTopicPart, "\n"); newlineIndex != -1 {
|
||||||
|
replyTopic = strings.TrimSpace(replyTopicPart[:newlineIndex])
|
||||||
|
content = strings.TrimSpace(content[:lastIndex]) + strings.TrimSpace(replyTopicPart[newlineIndex:])
|
||||||
|
} else {
|
||||||
|
replyTopic = strings.TrimSpace(replyTopicPart)
|
||||||
|
content = strings.TrimSpace(content[:lastIndex])
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -285,8 +305,15 @@ func (c *MQTTChannel) Send(ctx context.Context, msg bus.OutboundMessage) error {
|
||||||
// Replace placeholders in reply topic
|
// Replace placeholders in reply topic
|
||||||
replyTopic = strings.ReplaceAll(replyTopic, "{client_id}", c.config.ClientID)
|
replyTopic = strings.ReplaceAll(replyTopic, "{client_id}", c.config.ClientID)
|
||||||
replyTopic = strings.ReplaceAll(replyTopic, "{topic}", msg.ChatID)
|
replyTopic = strings.ReplaceAll(replyTopic, "{topic}", msg.ChatID)
|
||||||
|
replyTopic = strings.ReplaceAll(replyTopic, "{timestamp}", fmt.Sprintf("%d", time.Now().Unix()))
|
||||||
|
replyTopic = strings.ReplaceAll(replyTopic, "{message_id}", msg.ReplyToMessageID)
|
||||||
// Add more placeholders as needed
|
// Add more placeholders as needed
|
||||||
|
|
||||||
|
// Validate reply topic
|
||||||
|
if replyTopic == "" {
|
||||||
|
return fmt.Errorf("reply topic is empty and no default configured")
|
||||||
|
}
|
||||||
|
|
||||||
var payload []byte
|
var payload []byte
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue