mirror of
https://github.com/dpup/meshstream.git
synced 2026-08-07 09:12:55 +02:00
Updated logging configuration for dev and json for prod
This commit is contained in:
+1
-7
@@ -18,17 +18,11 @@ type Broker struct {
|
||||
|
||||
// NewBroker creates a new broker that distributes messages from sourceChannel to subscribers
|
||||
func NewBroker(sourceChannel <-chan *Packet, logger logging.Logger) *Broker {
|
||||
// Create a named logger if one was not provided
|
||||
if logger == nil {
|
||||
logger = logging.NewDevLogger()
|
||||
}
|
||||
brokerLogger := logger.Named("mqtt.broker")
|
||||
|
||||
broker := &Broker{
|
||||
sourceChan: sourceChannel,
|
||||
subscribers: make(map[chan *Packet]struct{}),
|
||||
done: make(chan struct{}),
|
||||
logger: brokerLogger,
|
||||
logger: logger.Named("mqtt.broker"),
|
||||
}
|
||||
|
||||
// Start the dispatch loop
|
||||
|
||||
+7
-13
@@ -30,17 +30,11 @@ type Client struct {
|
||||
|
||||
// NewClient creates a new MQTT client with the provided configuration
|
||||
func NewClient(config Config, logger logging.Logger) *Client {
|
||||
// Use provided logger or create a default one
|
||||
if logger == nil {
|
||||
logger = logging.NewDevLogger()
|
||||
}
|
||||
clientLogger := logger.Named("mqtt.client")
|
||||
|
||||
return &Client{
|
||||
config: config,
|
||||
decodedMessages: make(chan *Packet, 100), // Buffer up to 100 messages
|
||||
done: make(chan struct{}),
|
||||
logger: clientLogger,
|
||||
logger: logger.Named("mqtt.client"),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -94,7 +88,7 @@ func (c *Client) messageHandler(client mqtt.Client, msg mqtt.Message) {
|
||||
// Parse the topic structure
|
||||
topicInfo, err := decoder.ParseTopic(msg.Topic())
|
||||
if err != nil {
|
||||
c.logger.Errorw("Error parsing topic",
|
||||
c.logger.Errorw("Error parsing topic",
|
||||
"error", err,
|
||||
"topic", msg.Topic(),
|
||||
"payload_hex", fmt.Sprintf("%x", msg.Payload()),
|
||||
@@ -107,13 +101,13 @@ func (c *Client) messageHandler(client mqtt.Client, msg mqtt.Message) {
|
||||
case "e", "c", "map":
|
||||
// Binary encoded protobuf message
|
||||
decodedPacket := decoder.DecodeMessage(msg.Payload(), topicInfo)
|
||||
|
||||
|
||||
// Create packet with both the decoded packet and topic info
|
||||
packet := &Packet{
|
||||
DecodedPacket: decodedPacket,
|
||||
TopicInfo: topicInfo,
|
||||
}
|
||||
|
||||
|
||||
// Send the decoded message to the channel, but don't block if buffer is full
|
||||
select {
|
||||
case c.decodedMessages <- packet:
|
||||
@@ -125,11 +119,11 @@ func (c *Client) messageHandler(client mqtt.Client, msg mqtt.Message) {
|
||||
// Channel buffer is full, log a warning and drop the message
|
||||
c.logger.Warn("Message buffer full, dropping message")
|
||||
}
|
||||
|
||||
|
||||
case "json":
|
||||
// TODO: Add support for JSON format messages in the future
|
||||
c.logger.Debugf("Ignoring JSON format message from topic: %s", msg.Topic())
|
||||
|
||||
|
||||
default:
|
||||
// Unsupported format, log and ignore
|
||||
c.logger.Infow("Unsupported format", "format", topicInfo.Format, "topic", msg.Topic())
|
||||
@@ -144,4 +138,4 @@ func (c *Client) connectHandler(client mqtt.Client) {
|
||||
// connectionLostHandler is called when the client loses connection
|
||||
func (c *Client) connectionLostHandler(client mqtt.Client, err error) {
|
||||
c.logger.Errorw("Connection lost", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -42,6 +42,7 @@ func NewMessageStats(broker *Broker, printInterval time.Duration, logger logging
|
||||
Processor: s.recordMessage,
|
||||
StartHook: func() { go s.runTicker() },
|
||||
CloseHook: func() { s.ticker.Stop() },
|
||||
Logger: statsLogger,
|
||||
})
|
||||
|
||||
// Start processing messages
|
||||
|
||||
+2
-7
@@ -33,13 +33,8 @@ type BaseSubscriber struct {
|
||||
|
||||
// NewBaseSubscriber creates a new base subscriber
|
||||
func NewBaseSubscriber(config SubscriberConfig) *BaseSubscriber {
|
||||
// Use provided logger or create a default one
|
||||
var subscriberLogger logging.Logger
|
||||
if config.Logger == nil {
|
||||
subscriberLogger = logging.NewDevLogger().Named("mqtt.subscriber." + config.Name)
|
||||
} else {
|
||||
subscriberLogger = config.Logger.Named("mqtt.subscriber." + config.Name)
|
||||
}
|
||||
|
||||
subscriberLogger := config.Logger.Named("mqtt.subscriber." + config.Name)
|
||||
|
||||
return &BaseSubscriber{
|
||||
broker: config.Broker,
|
||||
|
||||
Reference in New Issue
Block a user