package rabbitmq import ( "context" "encoding/json" "fmt" "time" ) // ReadingEvent mirrors ingest-service's newReadingEvent wire format. type ReadingEvent struct { DeviceID string `json:"device_id"` ZoneID string `json:"zone_id"` SensorType string `json:"sensor_type"` Value float64 `json:"value"` RecordedAt time.Time `json:"recorded_at"` } // HandleResult tells ConsumeReadings how to settle a delivery. type HandleResult int const ( Ack HandleResult = iota // processed successfully NackRequeue // transient failure (e.g. gRPC unreachable) — try again later NackDiscard // permanent failure (e.g. bad data) — retrying won't help ) // Handler processes one reading event and decides how it should be settled. type Handler func(ctx context.Context, event ReadingEvent) HandleResult // ConsumeReadings blocks, delivering messages to handler, until ctx is // cancelled or the delivery channel closes (e.g. connection lost). func (c *Client) ConsumeReadings(ctx context.Context, handler Handler) error { deliveries, err := c.consumeCh.ConsumeWithContext(ctx, NewReadingQueue, "", false, false, false, false, nil) if err != nil { return fmt.Errorf("consume %s: %w", NewReadingQueue, err) } for { select { case <-ctx.Done(): return nil case d, ok := <-deliveries: if !ok { return fmt.Errorf("delivery channel for %s closed", NewReadingQueue) } var event ReadingEvent if err := json.Unmarshal(d.Body, &event); err != nil { // Poison message: no amount of retrying fixes malformed JSON. _ = d.Nack(false, false) continue } switch handler(ctx, event) { case Ack: _ = d.Ack(false) case NackRequeue: _ = d.Nack(false, true) case NackDiscard: _ = d.Nack(false, false) } } } }