// 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" "time" "github.com/rs/zerolog" 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/econetclient" "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 log *zerolog.Logger client *econetclient.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 e.log = zerolog.Ctx(ctx) // Fail fast on missing credentials (cheap, no network). if cfg.EconetEmail == "" || cfg.EconetPassword == "" { return nil, fmt.Errorf("econet: email and password required (set ECONET_EMAIL / ECONET_PASSWORD)") } e.client = econetclient.New(cfg, e.log) e.grpc = econetgrpc.NewEconetGRPCServer(ctx, cfg, e.client) e.mcp = econetmcp.NewEconetMCPServer(ctx, cfg, e.grpc) e.metrics = econetmetrics.NewCollector(ctx, cfg, e.client) if err := e.metrics.RegisterGauges(); err != nil { return nil, err } // Poll the cloud in the background so a slow or unreachable Rheem cloud // never blocks HTTP/gRPC server startup. Devices are empty until the first // successful refresh; the gauges and gRPC handlers tolerate that. go e.pollLoop(ctx) return e.shutdown, nil } // pollLoop refreshes device state on a ticker for the life of the service. // A per-refresh timeout keeps a hung REST call from stalling the loop. func (e *EconetService) pollLoop(ctx context.Context) { interval := e.config.GetPollInterval() e.refresh(ctx) t := time.NewTicker(interval) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: e.refresh(ctx) } } } func (e *EconetService) refresh(ctx context.Context) { rctx, cancel := context.WithTimeout(ctx, e.config.GetPollInterval()) defer cancel() if err := e.client.Refresh(rctx); err != nil { e.log.Error().Err(err).Msg("econet: refresh failed") return } e.log.Debug().Int("devices", len(e.client.Devices())).Msg("econet: refreshed") } func (e *EconetService) shutdown(_ context.Context) (string, error) { return "EconetService", nil } 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 yet") } return nil }, } }