generated from rmcguire/go-server-with-otel
145 lines
4.6 KiB
Go
145 lines
4.6 KiB
Go
// Package econetgrpc implements the EconetService gRPC API over a shared
|
|
// econetclient.Client, exposing read-only device state.
|
|
package econetgrpc
|
|
|
|
import (
|
|
"context"
|
|
|
|
"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"
|
|
"gitea.libretechconsulting.com/rmcguire/econet-exporter/pkg/econet/econetclient"
|
|
)
|
|
|
|
type EconetGRPCServer struct {
|
|
tracer trace.Tracer
|
|
ctx context.Context
|
|
cfg *config.ServiceConfig
|
|
client *econetclient.Client
|
|
pb.UnimplementedEconetServiceServer
|
|
}
|
|
|
|
func NewEconetGRPCServer(ctx context.Context, cfg *config.ServiceConfig, client *econetclient.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
|
|
}
|
|
|
|
// DeviceRawJSON returns the undecoded EcoNet equipment payload for a device,
|
|
// used by the web UI's debug view. ok is false if the serial is unknown or no
|
|
// raw payload was captured. It is a debug accessor, not part of the RPC surface.
|
|
func (s *EconetGRPCServer) DeviceRawJSON(serial string) ([]byte, bool) {
|
|
d := s.client.Device(serial)
|
|
if d == nil || len(d.Raw) == 0 {
|
|
return nil, false
|
|
}
|
|
return d.Raw, true
|
|
}
|
|
|
|
// modeSlugs maps the settable proto Mode enum to the normalized slugs the
|
|
// client (and Device.mode) use. MODE_UNSPECIFIED is intentionally absent;
|
|
// protovalidate rejects it before we get here.
|
|
var modeSlugs = map[pb.Mode]string{
|
|
pb.Mode_MODE_OFF: "off",
|
|
pb.Mode_MODE_ELECTRIC: "electric",
|
|
pb.Mode_MODE_ENERGY_SAVING: "energy-saving",
|
|
pb.Mode_MODE_HEAT_PUMP: "heat-pump",
|
|
pb.Mode_MODE_HIGH_DEMAND: "high-demand",
|
|
pb.Mode_MODE_GAS: "gas",
|
|
pb.Mode_MODE_PERFORMANCE: "performance",
|
|
pb.Mode_MODE_VACATION: "vacation",
|
|
pb.Mode_MODE_ELECTRIC_GAS: "electric-gas",
|
|
}
|
|
|
|
// ModeSlug returns the normalized slug for a settable Mode, or "" if unknown.
|
|
func ModeSlug(m pb.Mode) string { return modeSlugs[m] }
|
|
|
|
// modeFromSlug is the reverse of modeSlugs; ok is false for slugs with no
|
|
// settable enum (e.g. "unknown").
|
|
func modeFromSlug(slug string) (pb.Mode, bool) {
|
|
for m, s := range modeSlugs {
|
|
if s == slug {
|
|
return m, true
|
|
}
|
|
}
|
|
return pb.Mode_MODE_UNSPECIFIED, false
|
|
}
|
|
|
|
func (s *EconetGRPCServer) SetMode(ctx context.Context, req *pb.SetModeRequest) (
|
|
*pb.SetModeResponse, error,
|
|
) {
|
|
ctx, span := s.tracer.Start(ctx, "setMode", trace.WithAttributes(
|
|
attribute.String("serialNumber", req.GetSerialNumber()),
|
|
attribute.String("mode", req.GetMode().String()),
|
|
))
|
|
defer span.End()
|
|
|
|
slug, ok := modeSlugs[req.GetMode()]
|
|
if !ok {
|
|
err := status.Errorf(grpccodes.InvalidArgument, "unsupported mode %q", req.GetMode())
|
|
span.SetStatus(codes.Error, err.Error())
|
|
return nil, err
|
|
}
|
|
|
|
if err := s.client.SetMode(ctx, req.GetSerialNumber(), slug); err != nil {
|
|
err := status.Error(grpccodes.FailedPrecondition, err.Error())
|
|
span.SetStatus(codes.Error, err.Error())
|
|
return nil, err
|
|
}
|
|
|
|
span.SetStatus(codes.Ok, "")
|
|
// State applies asynchronously; return the last-known snapshot. Guard against
|
|
// a concurrent refresh evicting the device between publish and lookup.
|
|
resp := &pb.SetModeResponse{}
|
|
if d := s.client.Device(req.GetSerialNumber()); d != nil {
|
|
resp.Device = deviceToProto(d)
|
|
}
|
|
return resp, nil
|
|
}
|