generated from rmcguire/go-server-with-otel
implement econet exporter
This commit is contained in:
@@ -0,0 +1,94 @@
|
||||
// Package econet provides EconetService, a service.AppService that connects
|
||||
// to the Rheem EcoNet cloud and exposes device state via gRPC, an MCP tool,
|
||||
// and OTEL metrics served for Prometheus scraping.
|
||||
package econet
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
rheemcloud "github.com/kevinburke/rheemcloud-go"
|
||||
|
||||
optsgrpc "gitea.libretechconsulting.com/rmcguire/go-app/pkg/srv/grpc/opts"
|
||||
optshttp "gitea.libretechconsulting.com/rmcguire/go-app/pkg/srv/http/opts"
|
||||
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/config"
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/econet/econetgrpc"
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/econet/econetmcp"
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/econet/econetmetrics"
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/service"
|
||||
)
|
||||
|
||||
type EconetService struct {
|
||||
ctx context.Context
|
||||
config *config.ServiceConfig
|
||||
client *rheemcloud.Client
|
||||
grpc *econetgrpc.EconetGRPCServer
|
||||
mcp *econetmcp.EconetMCPServer
|
||||
metrics *econetmetrics.Collector
|
||||
}
|
||||
|
||||
func (e *EconetService) Init(ctx context.Context, cfg *config.ServiceConfig) (service.ShutdownFunc, error) {
|
||||
e.ctx = ctx
|
||||
e.config = cfg
|
||||
|
||||
client, err := connect(ctx, cfg)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
e.client = client
|
||||
|
||||
e.grpc = econetgrpc.NewEconetGRPCServer(ctx, cfg, client)
|
||||
e.mcp = econetmcp.NewEconetMCPServer(ctx, cfg, e.grpc)
|
||||
e.metrics = econetmetrics.NewCollector(ctx, cfg, client)
|
||||
if err := e.metrics.Start(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return e.shutdown, nil
|
||||
}
|
||||
|
||||
func connect(ctx context.Context, cfg *config.ServiceConfig) (*rheemcloud.Client, error) {
|
||||
if cfg.EconetEmail == "" || cfg.EconetPassword == "" {
|
||||
return nil, fmt.Errorf("econet: email and password required (set ECONET_EMAIL / ECONET_PASSWORD)")
|
||||
}
|
||||
client, err := rheemcloud.Connect(ctx, cfg.EconetEmail, cfg.EconetPassword, &rheemcloud.Config{
|
||||
Logger: slog.Default(),
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("econet: connect failed: %w", err)
|
||||
}
|
||||
return client, nil
|
||||
}
|
||||
|
||||
func (e *EconetService) shutdown(_ context.Context) (string, error) {
|
||||
e.metrics.Stop()
|
||||
return "EconetService", e.client.Close()
|
||||
}
|
||||
|
||||
func (e *EconetService) GetGRPC() *optsgrpc.AppGRPC {
|
||||
return &optsgrpc.AppGRPC{
|
||||
Services: e.grpc.GetServices(),
|
||||
GRPCDialOpts: e.grpc.GetDialOpts(),
|
||||
}
|
||||
}
|
||||
|
||||
func (e *EconetService) GetHTTP() *optshttp.AppHTTP {
|
||||
return &optshttp.AppHTTP{
|
||||
Ctx: e.ctx,
|
||||
Handlers: e.mcp.GetHandlers(),
|
||||
HealthChecks: e.healthChecks(),
|
||||
}
|
||||
}
|
||||
|
||||
func (e *EconetService) healthChecks() []optshttp.HealthCheckFunc {
|
||||
return []optshttp.HealthCheckFunc{
|
||||
func(_ context.Context) error {
|
||||
if len(e.client.Devices()) == 0 {
|
||||
return fmt.Errorf("econet: no devices loaded")
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package econetgrpc
|
||||
|
||||
import (
|
||||
rheemcloud "github.com/kevinburke/rheemcloud-go"
|
||||
"google.golang.org/protobuf/types/known/timestamppb"
|
||||
|
||||
pb "gitea.libretechconsulting.com/rmcguire/econet-exporter/api/econet/v1alpha1"
|
||||
)
|
||||
|
||||
func deviceToProto(d *rheemcloud.Device) *pb.Device {
|
||||
min, max := d.SetpointLimits()
|
||||
return &pb.Device{
|
||||
SerialNumber: d.SerialNumber(),
|
||||
DeviceId: d.DeviceID(),
|
||||
FriendlyName: d.FriendlyName(),
|
||||
Type: d.Type().String(),
|
||||
GenericType: d.GenericType(),
|
||||
Connected: d.Connected(),
|
||||
WifiSignal: int32(d.WiFiSignal()),
|
||||
Mode: d.Mode().String(),
|
||||
Enabled: d.Enabled(),
|
||||
Running: d.Running(),
|
||||
RunningState: d.RunningState(),
|
||||
Setpoint: int32(d.Setpoint()),
|
||||
SetpointMin: int32(min),
|
||||
SetpointMax: int32(max),
|
||||
HotWaterAvailability: int32(d.HotWaterAvailability()),
|
||||
AlertCount: int32(d.AlertCount()),
|
||||
Away: d.Away(),
|
||||
LastUpdated: timestamppb.New(d.LastUpdated),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
// Package econetgrpc implements the EconetService gRPC API over a shared
|
||||
// rheemcloud.Client, exposing read-only device state.
|
||||
package econetgrpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
rheemcloud "github.com/kevinburke/rheemcloud-go"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
grpccodes "google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
|
||||
"gitea.libretechconsulting.com/rmcguire/go-app/pkg/otel"
|
||||
|
||||
pb "gitea.libretechconsulting.com/rmcguire/econet-exporter/api/econet/v1alpha1"
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/config"
|
||||
)
|
||||
|
||||
type EconetGRPCServer struct {
|
||||
tracer trace.Tracer
|
||||
ctx context.Context
|
||||
cfg *config.ServiceConfig
|
||||
client *rheemcloud.Client
|
||||
pb.UnimplementedEconetServiceServer
|
||||
}
|
||||
|
||||
func NewEconetGRPCServer(ctx context.Context, cfg *config.ServiceConfig, client *rheemcloud.Client) *EconetGRPCServer {
|
||||
return &EconetGRPCServer{
|
||||
ctx: ctx,
|
||||
cfg: cfg,
|
||||
client: client,
|
||||
tracer: otel.GetTracer(ctx, "econetGRPCServer"),
|
||||
}
|
||||
}
|
||||
|
||||
func (s *EconetGRPCServer) ListDevices(ctx context.Context, _ *pb.ListDevicesRequest) (
|
||||
*pb.ListDevicesResponse, error,
|
||||
) {
|
||||
_, span := s.tracer.Start(ctx, "listDevices")
|
||||
defer span.End()
|
||||
|
||||
devices := s.client.Devices()
|
||||
resp := &pb.ListDevicesResponse{Devices: make([]*pb.Device, 0, len(devices))}
|
||||
for _, d := range devices {
|
||||
resp.Devices = append(resp.Devices, deviceToProto(d))
|
||||
}
|
||||
|
||||
span.SetAttributes(attribute.Int("devices", len(resp.Devices)))
|
||||
span.SetStatus(codes.Ok, "")
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *EconetGRPCServer) GetDevice(ctx context.Context, req *pb.GetDeviceRequest) (
|
||||
*pb.GetDeviceResponse, error,
|
||||
) {
|
||||
_, span := s.tracer.Start(ctx, "getDevice", trace.WithAttributes(
|
||||
attribute.String("serialNumber", req.GetSerialNumber()),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
d := s.client.Device(req.GetSerialNumber())
|
||||
if d == nil {
|
||||
err := status.Errorf(grpccodes.NotFound, "no device with serial %q", req.GetSerialNumber())
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return nil, err
|
||||
}
|
||||
|
||||
span.SetStatus(codes.Ok, "")
|
||||
return &pb.GetDeviceResponse{Device: deviceToProto(d)}, nil
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
package econetgrpc
|
||||
|
||||
import (
|
||||
"gitea.libretechconsulting.com/rmcguire/go-app/pkg/srv/grpc/opts"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
|
||||
econetPb "gitea.libretechconsulting.com/rmcguire/econet-exporter/api/econet/v1alpha1"
|
||||
)
|
||||
|
||||
func (s *EconetGRPCServer) GetDialOpts() []grpc.DialOption {
|
||||
return []grpc.DialOption{
|
||||
// NOTE: Necessary for grpc-gateway to connect to grpc server
|
||||
// Update if grpc service has credentials, tls, etc..
|
||||
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
||||
}
|
||||
}
|
||||
|
||||
func (s *EconetGRPCServer) GetServices() []*opts.GRPCService {
|
||||
return []*opts.GRPCService{
|
||||
{
|
||||
Name: "Econet Exporter",
|
||||
Type: &econetPb.EconetService_ServiceDesc,
|
||||
Service: s,
|
||||
GwRegistrationFuncs: []opts.GwRegistrationFunc{
|
||||
econetPb.RegisterEconetServiceHandler,
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
/*
|
||||
Package econetmcp exposes an MCP server, mounted as an HTTP handler at
|
||||
/api/mcp, whose tools wrap the EconetService gRPC API.
|
||||
*/
|
||||
package econetmcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
|
||||
"gitea.libretechconsulting.com/rmcguire/go-app/pkg/srv/http/opts"
|
||||
"github.com/modelcontextprotocol/go-sdk/mcp"
|
||||
"github.com/rs/zerolog"
|
||||
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/config"
|
||||
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/econet/econetgrpc"
|
||||
)
|
||||
|
||||
var EconetMCPImpl = &mcp.Implementation{
|
||||
Name: "Econet MCP Server",
|
||||
Title: "Econet Exporter MCP",
|
||||
}
|
||||
|
||||
type EconetMCPServer struct {
|
||||
ctx context.Context
|
||||
cfg *config.ServiceConfig
|
||||
log *zerolog.Logger
|
||||
server *mcp.Server
|
||||
econetGRPC *econetgrpc.EconetGRPCServer
|
||||
}
|
||||
|
||||
func NewEconetMCPServer(ctx context.Context, cfg *config.ServiceConfig,
|
||||
grpc *econetgrpc.EconetGRPCServer,
|
||||
) *EconetMCPServer {
|
||||
return &EconetMCPServer{
|
||||
ctx: ctx,
|
||||
cfg: cfg,
|
||||
log: zerolog.Ctx(ctx),
|
||||
econetGRPC: grpc,
|
||||
server: mcp.NewServer(EconetMCPImpl, &mcp.ServerOptions{
|
||||
Instructions: "Use these tools to inspect Rheem EcoNet water heaters: " +
|
||||
"current mode, setpoint, hot water availability, connectivity, and alerts.",
|
||||
HasTools: true,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
func (m *EconetMCPServer) GetHandlers() []opts.HTTPHandler {
|
||||
// NOTE: Add other tools here
|
||||
m.addListDevicesTool()
|
||||
|
||||
handler := mcp.NewStreamableHTTPHandler(func(*http.Request) *mcp.Server {
|
||||
return m.server
|
||||
}, &mcp.StreamableHTTPOptions{})
|
||||
|
||||
m.log.Debug().Msg("Econet MCP tools ready")
|
||||
|
||||
return []opts.HTTPHandler{
|
||||
{
|
||||
Prefix: "/api/mcp",
|
||||
Handler: handler,
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
package econetmcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/modelcontextprotocol/go-sdk/mcp"
|
||||
"k8s.io/utils/ptr"
|
||||
|
||||
pb "gitea.libretechconsulting.com/rmcguire/econet-exporter/api/econet/v1alpha1"
|
||||
)
|
||||
|
||||
var ListDevicesTool = &mcp.Tool{
|
||||
Name: "list_water_heaters",
|
||||
Title: "List EcoNet Water Heaters",
|
||||
Description: "Lists Rheem EcoNet devices and their current state (mode, setpoint, " +
|
||||
"hot water availability, connectivity). Optionally filter to one serial number.",
|
||||
Annotations: &mcp.ToolAnnotations{
|
||||
ReadOnlyHint: true,
|
||||
OpenWorldHint: ptr.To(true),
|
||||
},
|
||||
}
|
||||
|
||||
type ListDevicesParams struct {
|
||||
SerialNumber string `json:"serial_number,omitempty" jsonschema:"Optional device serial number to return a single device"`
|
||||
}
|
||||
|
||||
func (m *EconetMCPServer) addListDevicesTool() {
|
||||
mcp.AddTool(m.server, ListDevicesTool, m.listDevicesHandler)
|
||||
}
|
||||
|
||||
func (m *EconetMCPServer) listDevicesHandler(ctx context.Context, _ *mcp.CallToolRequest, args *ListDevicesParams) (
|
||||
*mcp.CallToolResult, any, error,
|
||||
) {
|
||||
m.log.Debug().Str("serial", args.SerialNumber).Msg("list water heaters tool called")
|
||||
|
||||
devices, err := m.lookup(ctx, args.SerialNumber)
|
||||
if err != nil {
|
||||
return &mcp.CallToolResult{IsError: true, Meta: mcp.Meta{"error": err.Error()}}, nil, err
|
||||
}
|
||||
|
||||
return &mcp.CallToolResult{
|
||||
Content: []mcp.Content{&mcp.TextContent{Text: summarizeDevices(devices)}},
|
||||
}, nil, nil
|
||||
}
|
||||
|
||||
func (m *EconetMCPServer) lookup(ctx context.Context, serial string) ([]*pb.Device, error) {
|
||||
if serial != "" {
|
||||
resp, err := m.econetGRPC.GetDevice(ctx, &pb.GetDeviceRequest{SerialNumber: serial})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return []*pb.Device{resp.GetDevice()}, nil
|
||||
}
|
||||
|
||||
resp, err := m.econetGRPC.ListDevices(ctx, &pb.ListDevicesRequest{})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return resp.GetDevices(), nil
|
||||
}
|
||||
|
||||
func summarizeDevices(devices []*pb.Device) string {
|
||||
if len(devices) == 0 {
|
||||
return "No EcoNet devices found."
|
||||
}
|
||||
var b strings.Builder
|
||||
for _, d := range devices {
|
||||
fmt.Fprintf(&b, "%s (%s, serial %s): mode=%s setpoint=%d°F running=%t connected=%t hotWater=%d%% alerts=%d\n",
|
||||
d.GetFriendlyName(), d.GetGenericType(), d.GetSerialNumber(),
|
||||
d.GetMode(), d.GetSetpoint(), d.GetRunning(), d.GetConnected(),
|
||||
d.GetHotWaterAvailability(), d.GetAlertCount())
|
||||
}
|
||||
return strings.TrimRight(b.String(), "\n")
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
package econetmcp
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
pb "gitea.libretechconsulting.com/rmcguire/econet-exporter/api/econet/v1alpha1"
|
||||
)
|
||||
|
||||
func TestSummarizeDevices(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
devices []*pb.Device
|
||||
want string // substring that must appear
|
||||
}{
|
||||
{"empty", nil, "No EcoNet devices found."},
|
||||
{
|
||||
name: "single",
|
||||
devices: []*pb.Device{{
|
||||
FriendlyName: "Garage",
|
||||
GenericType: "heatpumpWaterHeater",
|
||||
SerialNumber: "ABC123",
|
||||
Mode: "heat-pump",
|
||||
Setpoint: 125,
|
||||
}},
|
||||
want: "Garage (heatpumpWaterHeater, serial ABC123): mode=heat-pump setpoint=125",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
got := summarizeDevices(tt.devices)
|
||||
if !strings.Contains(got, tt.want) {
|
||||
t.Errorf("summarizeDevices() = %q, want substring %q", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package econetmetrics
|
||||
|
||||
import "go.opentelemetry.io/otel/metric"
|
||||
|
||||
// gaugeBuilder creates observable gauges, accumulating them for a single
|
||||
// RegisterCallback and short-circuiting on the first creation error.
|
||||
//
|
||||
// Units are intentionally baked into the metric names (Prometheus idiom)
|
||||
// rather than set via metric.WithUnit — the OTEL→Prometheus exporter would
|
||||
// otherwise append a second unit suffix (e.g. econet_setpoint_fahrenheit_F).
|
||||
type gaugeBuilder struct {
|
||||
m metric.Meter
|
||||
obs []metric.Observable
|
||||
err error
|
||||
}
|
||||
|
||||
func (b *gaugeBuilder) i64(name, desc string) metric.Int64ObservableGauge {
|
||||
if b.err != nil {
|
||||
return nil
|
||||
}
|
||||
g, err := b.m.Int64ObservableGauge(name, metric.WithDescription(desc))
|
||||
if err != nil {
|
||||
b.err = err
|
||||
return nil
|
||||
}
|
||||
b.obs = append(b.obs, g)
|
||||
return g
|
||||
}
|
||||
|
||||
func (b *gaugeBuilder) f64(name, desc string) metric.Float64ObservableGauge {
|
||||
if b.err != nil {
|
||||
return nil
|
||||
}
|
||||
g, err := b.m.Float64ObservableGauge(name, metric.WithDescription(desc))
|
||||
if err != nil {
|
||||
b.err = err
|
||||
return nil
|
||||
}
|
||||
b.obs = append(b.obs, g)
|
||||
return g
|
||||
}
|
||||
@@ -0,0 +1,132 @@
|
||||
package econetmetrics
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
rheemcloud "github.com/kevinburke/rheemcloud-go"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/metric"
|
||||
)
|
||||
|
||||
// registerLive wires the per-device live-state gauges to a single callback
|
||||
// that snapshots client.Devices() once per scrape (cheap, in-memory).
|
||||
func (c *Collector) registerLive() error {
|
||||
b := &gaugeBuilder{m: c.meter}
|
||||
g := &liveGauges{
|
||||
setpoint: b.i64("econet_setpoint_fahrenheit", "Target water temperature (°F)"),
|
||||
setpointMin: b.i64("econet_setpoint_min_fahrenheit", "Minimum allowed setpoint (°F)"),
|
||||
setpointMax: b.i64("econet_setpoint_max_fahrenheit", "Maximum allowed setpoint (°F)"),
|
||||
connected: b.i64("econet_connected", "1 if the device is connected"),
|
||||
running: b.i64("econet_running", "1 if the device is actively heating"),
|
||||
enabled: b.i64("econet_enabled", "1 if the device is enabled"),
|
||||
away: b.i64("econet_away", "1 if away mode is on"),
|
||||
hotWater: b.i64("econet_hot_water_availability_percent", "Hot water availability (%)"),
|
||||
alerts: b.i64("econet_alert_count", "Number of active alerts"),
|
||||
wifi: b.i64("econet_wifi_signal_db", "WiFi signal strength (dB)"),
|
||||
lastUpdated: b.i64("econet_last_updated_seconds", "Unix time of last MQTT state update"),
|
||||
info: b.i64("econet_device_info", "Device metadata (value is always 1)"),
|
||||
}
|
||||
if b.err != nil {
|
||||
return b.err
|
||||
}
|
||||
_, err := c.meter.RegisterCallback(func(_ context.Context, o metric.Observer) error {
|
||||
for _, d := range c.client.Devices() {
|
||||
g.observe(o, d)
|
||||
}
|
||||
return nil
|
||||
}, b.obs...)
|
||||
return err
|
||||
}
|
||||
|
||||
// registerUsage wires the polled energy/water gauges, reading the cache the
|
||||
// poller refreshes (the underlying REST calls are too slow for a callback).
|
||||
func (c *Collector) registerUsage() error {
|
||||
b := &gaugeBuilder{m: c.meter}
|
||||
g := &usageGauges{
|
||||
energy: b.f64("econet_energy_usage_kwh", "Energy used today (kWh)"),
|
||||
cost: b.f64("econet_energy_cost_dollars", "Energy cost today in US dollars (kWh * costPerKWH)"),
|
||||
water: b.f64("econet_water_usage_gallons", "Water used today (gallons)"),
|
||||
}
|
||||
if b.err != nil {
|
||||
return b.err
|
||||
}
|
||||
_, err := c.meter.RegisterCallback(func(_ context.Context, o metric.Observer) error {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
for serial, u := range c.usage {
|
||||
g.observe(o, serial, u, c.cfg.CostPerKWH)
|
||||
}
|
||||
return nil
|
||||
}, b.obs...)
|
||||
return err
|
||||
}
|
||||
|
||||
type liveGauges struct {
|
||||
setpoint, setpointMin, setpointMax metric.Int64ObservableGauge
|
||||
connected, running, enabled, away metric.Int64ObservableGauge
|
||||
hotWater, alerts, wifi, lastUpdated metric.Int64ObservableGauge
|
||||
info metric.Int64ObservableGauge
|
||||
}
|
||||
|
||||
func (g *liveGauges) observe(o metric.Observer, d *rheemcloud.Device) {
|
||||
set := metric.WithAttributes(deviceAttrs(d)...)
|
||||
min, max := d.SetpointLimits()
|
||||
o.ObserveInt64(g.setpoint, int64(d.Setpoint()), set)
|
||||
o.ObserveInt64(g.setpointMin, int64(min), set)
|
||||
o.ObserveInt64(g.setpointMax, int64(max), set)
|
||||
o.ObserveInt64(g.connected, b2i(d.Connected()), set)
|
||||
o.ObserveInt64(g.running, b2i(d.Running()), set)
|
||||
o.ObserveInt64(g.enabled, b2i(d.Enabled()), set)
|
||||
o.ObserveInt64(g.away, b2i(d.Away()), set)
|
||||
o.ObserveInt64(g.hotWater, int64(d.HotWaterAvailability()), set)
|
||||
o.ObserveInt64(g.alerts, int64(d.AlertCount()), set)
|
||||
o.ObserveInt64(g.wifi, int64(d.WiFiSignal()), set)
|
||||
// LastUpdated is zero until the first MQTT update lands; skip it then
|
||||
// rather than emit a nonsensical negative epoch.
|
||||
if !d.LastUpdated.IsZero() {
|
||||
o.ObserveInt64(g.lastUpdated, d.LastUpdated.Unix(), set)
|
||||
}
|
||||
o.ObserveInt64(g.info, 1, metric.WithAttributes(infoAttrs(d)...))
|
||||
}
|
||||
|
||||
type usageGauges struct {
|
||||
energy, cost, water metric.Float64ObservableGauge
|
||||
}
|
||||
|
||||
func (g *usageGauges) observe(o metric.Observer, serial string, u usage, costPerKWH float64) {
|
||||
set := metric.WithAttributes(attribute.String("econet.serial", serial))
|
||||
o.ObserveFloat64(g.energy, u.energyKWH, metric.WithAttributes(
|
||||
attribute.String("econet.serial", serial),
|
||||
attribute.String("econet.energy_type", u.energyType),
|
||||
))
|
||||
o.ObserveFloat64(g.water, u.waterGallons, set)
|
||||
if u.energyType == "KWH" {
|
||||
o.ObserveFloat64(g.cost, u.energyKWH*costPerKWH, set)
|
||||
}
|
||||
}
|
||||
|
||||
// deviceAttrs are low-cardinality identity attributes safe for every series.
|
||||
// They follow OTEL semantic-convention style (namespaced, dotted keys) and
|
||||
// deliberately exclude anything that changes over time (e.g. timestamps).
|
||||
func deviceAttrs(d *rheemcloud.Device) []attribute.KeyValue {
|
||||
return []attribute.KeyValue{
|
||||
attribute.String("econet.serial", d.SerialNumber()),
|
||||
attribute.String("econet.device_id", d.DeviceID()),
|
||||
attribute.String("econet.friendly_name", d.FriendlyName()),
|
||||
}
|
||||
}
|
||||
|
||||
func infoAttrs(d *rheemcloud.Device) []attribute.KeyValue {
|
||||
return append(deviceAttrs(d),
|
||||
attribute.String("econet.type", d.Type().String()),
|
||||
attribute.String("econet.generic_type", d.GenericType()),
|
||||
attribute.String("econet.mode", d.Mode().String()),
|
||||
)
|
||||
}
|
||||
|
||||
func b2i(b bool) int64 {
|
||||
if b {
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package econetmetrics
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestB2I(t *testing.T) {
|
||||
if got := b2i(true); got != 1 {
|
||||
t.Errorf("b2i(true) = %d, want 1", got)
|
||||
}
|
||||
if got := b2i(false); got != 0 {
|
||||
t.Errorf("b2i(false) = %d, want 0", got)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,123 @@
|
||||
// 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()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user