// Package econetmetrics registers OTEL observable gauges for Rheem EcoNet // device state and periodically polls energy/water usage history so the // values are available at Prometheus scrape time. package econetmetrics import ( "context" "sync" "time" rheemcloud "github.com/kevinburke/rheemcloud-go" "github.com/rs/zerolog" "go.opentelemetry.io/otel/metric" "gitea.libretechconsulting.com/rmcguire/go-app/pkg/otel" "gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/config" ) // usage holds the most recent polled energy/water totals for a device. type usage struct { energyKWH float64 energyType string waterGallons float64 } type Collector struct { ctx context.Context cfg *config.ServiceConfig log *zerolog.Logger client *rheemcloud.Client meter metric.Meter mu sync.RWMutex usage map[string]usage // keyed by serial number stop chan struct{} } func NewCollector(ctx context.Context, cfg *config.ServiceConfig, client *rheemcloud.Client) *Collector { return &Collector{ ctx: ctx, cfg: cfg, log: zerolog.Ctx(ctx), client: client, meter: otel.GetMeter(ctx, "econet"), usage: make(map[string]usage), stop: make(chan struct{}), } } // Start registers the gauge callbacks and launches the usage poller and // event drain goroutines. The rheemcloud client keeps live state fresh over // MQTT; draining its event channel keeps the cap-64 buffer from filling. func (c *Collector) Start() error { if err := c.registerLive(); err != nil { return err } if err := c.registerUsage(); err != nil { return err } go c.drainEvents() go c.pollUsage() return nil } func (c *Collector) Stop() { close(c.stop) } func (c *Collector) drainEvents() { events := c.client.Subscribe() for { select { case <-c.stop: return case <-c.ctx.Done(): return case ev, ok := <-events: if !ok { return } c.log.Debug().Str("event", ev.Kind.String()).Msg("econet event") } } } func (c *Collector) pollUsage() { c.refreshUsage() t := time.NewTicker(c.cfg.GetUsageInterval()) defer t.Stop() for { select { case <-c.stop: return case <-c.ctx.Done(): return case <-t.C: c.refreshUsage() } } } // refreshUsage polls each device's energy/water usage for the current day // and caches the totals for the usage gauge callback to read. func (c *Collector) refreshUsage() { now := time.Now() start := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, now.Location()) for serial, d := range c.client.Devices() { u := usage{} if e, err := c.client.EnergyUsage(c.ctx, d, start, now); err == nil { u.energyKWH, u.energyType = e.Total, e.EnergyType } else { c.log.Debug().Err(err).Str("serial", serial).Msg("energy usage poll failed") } if w, err := c.client.WaterUsage(c.ctx, d, start, now); err == nil { u.waterGallons = w.Total } else { c.log.Debug().Err(err).Str("serial", serial).Msg("water usage poll failed") } c.mu.Lock() c.usage[serial] = u c.mu.Unlock() } }