moved all validation and parsing from handler to service layer

This commit is contained in:
Pedro Pérez 2025-10-10 02:42:29 +02:00
parent 8774b55d3d
commit 40ffee4d56
4 changed files with 67 additions and 53 deletions

View File

@ -106,6 +106,8 @@ func (r *SensorDataRequest) Validate() error {
} }
var ( var (
ErrRegisteringSensor = errors.New("error registering sensor")
ErrUpdatingSensor = errors.New("error updating sensor")
ErrInvalidSensorIdentifier = errors.New("sensor identifier is required") ErrInvalidSensorIdentifier = errors.New("sensor identifier is required")
ErrInvalidSensorType = errors.New("sensor type is required") ErrInvalidSensorType = errors.New("sensor type is required")
ErrSensorNotFound = errors.New("sensor not found") ErrSensorNotFound = errors.New("sensor not found")

View File

@ -4,7 +4,6 @@ import (
"encoding/json" "encoding/json"
"log/slog" "log/slog"
"nats-app/internal/iot" "nats-app/internal/iot"
"time"
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
) )
@ -93,11 +92,6 @@ func (h *Handlers) SetupEndpoints() *Handlers {
func (h *Handlers) register() { func (h *Handlers) register() {
h.NATS.Subscribe(subjectSensorsRegister, func(msg *nats.Msg) { h.NATS.Subscribe(subjectSensorsRegister, func(msg *nats.Msg) {
handleRequest(msg, func(req Sensor) (Sensor, error) { handleRequest(msg, func(req Sensor) (Sensor, error) {
if err := req.Validate(); err != nil {
slog.Error("error validating sensor", "error", err)
return Sensor{}, err
}
if err := h.service.RegisterSensor(req); err != nil { if err := h.service.RegisterSensor(req); err != nil {
return Sensor{}, err return Sensor{}, err
} }
@ -112,18 +106,7 @@ func (h *Handlers) register() {
func (h *Handlers) registerData() { func (h *Handlers) registerData() {
h.NATS.Subscribe(subjectSensorsData+"*", func(msg *nats.Msg) { h.NATS.Subscribe(subjectSensorsData+"*", func(msg *nats.Msg) {
handlePublish(msg, func(data SensorData) error { handlePublish(msg, func(data SensorData) error {
if err := data.Validate(); err != nil { return h.service.RegisterSensorData(data)
slog.Error("error validating sensor data", "error", err)
return err
}
if err := h.service.RegisterSensorData(data); err != nil {
slog.Error("failed to save sensor data", "error", err, "sensor_id", data.SensorID)
return err
}
slog.Debug("sensor data saved", "sensor_id", data.SensorID, "value", data.Value)
return nil
}) })
}) })
} }
@ -131,12 +114,6 @@ func (h *Handlers) registerData() {
func (h *Handlers) update() { func (h *Handlers) update() {
h.NATS.Subscribe(subjectSensorsUpdate, func(msg *nats.Msg) { h.NATS.Subscribe(subjectSensorsUpdate, func(msg *nats.Msg) {
handleRequest(msg, func(req Sensor) (Sensor, error) { handleRequest(msg, func(req Sensor) (Sensor, error) {
slog.Debug("calling sensor.update", "payload", req)
if err := req.Validate(); err != nil {
return Sensor{}, err
}
if err := h.service.UpdateSensor(req); err != nil { if err := h.service.UpdateSensor(req); err != nil {
return Sensor{}, err return Sensor{}, err
} }
@ -151,13 +128,7 @@ func (h *Handlers) update() {
func (h *Handlers) get() { func (h *Handlers) get() {
h.NATS.Subscribe(subjectSensorsGet, func(msg *nats.Msg) { h.NATS.Subscribe(subjectSensorsGet, func(msg *nats.Msg) {
handleRequest(msg, func(req SensorRequest) (Sensor, error) { handleRequest(msg, func(req SensorRequest) (Sensor, error) {
slog.Debug("calling sensor.get", "payload", req) return h.service.GetSensor(req)
if err := req.Validate(); err != nil {
return Sensor{}, err
}
return h.service.GetSensor(req.SensorID)
}) })
}) })
} }
@ -165,21 +136,7 @@ func (h *Handlers) get() {
func (h *Handlers) getValues() { func (h *Handlers) getValues() {
h.NATS.Subscribe(subjectSensorsValuesGet, func(msg *nats.Msg) { h.NATS.Subscribe(subjectSensorsValuesGet, func(msg *nats.Msg) {
handleRequest(msg, func(req SensorDataRequest) ([]SensorData, error) { handleRequest(msg, func(req SensorDataRequest) ([]SensorData, error) {
if err := req.Validate(); err != nil { return h.service.GetValues(req)
return []SensorData{}, err
}
from, err := time.Parse(time.RFC3339, *req.From)
if err != nil {
return []SensorData{}, err
}
to, err := time.Parse(time.RFC3339, *req.To)
if err != nil {
return []SensorData{}, err
}
return h.service.GetValues(req.SensorID, from, to)
}) })
}) })
} }

View File

@ -2,6 +2,7 @@ package sensors
import ( import (
"log/slog" "log/slog"
"strings"
"time" "time"
) )
@ -16,20 +17,32 @@ func NewService(repo Repository) *Service {
} }
func (s *Service) RegisterSensor(sensor Sensor) error { func (s *Service) RegisterSensor(sensor Sensor) error {
if err := sensor.Validate(); err != nil {
slog.Error("error validating sensor", "error", err)
return err
}
err := s.repo.CreateSensor(sensor) err := s.repo.CreateSensor(sensor)
if err != nil { if err != nil {
slog.Error("error registering sensor", "error", err) slog.Error("error registering sensor", "error", err)
return err if strings.Contains(err.Error(), "duplicate key value") {
return ErrSensorAlreadyExists
}
return ErrRegisteringSensor
} }
return nil return nil
} }
func (s *Service) RegisterSensorData(data SensorData) error { func (s *Service) RegisterSensorData(data SensorData) error {
if err := data.Validate(); err != nil {
slog.Error("error validating sensor data", "error", err)
return err
}
err := s.repo.CreateSensorData(data) err := s.repo.CreateSensorData(data)
if err != nil { if err != nil {
slog.Error("error registering sensor data") slog.Error("error registering sensor data", "error", err)
return err return err
} }
@ -37,15 +50,51 @@ func (s *Service) RegisterSensorData(data SensorData) error {
} }
func (s *Service) UpdateSensor(sensor Sensor) error { func (s *Service) UpdateSensor(sensor Sensor) error {
return s.repo.UpdateSensor(sensor) if err := sensor.Validate(); err != nil {
slog.Error("error validating sensor data", "error", err)
return err
}
err := s.repo.UpdateSensor(sensor)
if err != nil {
slog.Error("error updating sensor", "error", err)
if strings.Contains(err.Error(), "duplicate key value") {
return ErrSensorAlreadyExists
}
return ErrUpdatingSensor
}
return nil
} }
func (s *Service) GetSensor(sensorID string) (Sensor, error) { func (s *Service) GetSensor(sensor SensorRequest) (Sensor, error) {
return s.repo.ReadSensor(sensorID) if err := sensor.Validate(); err != nil {
slog.Error("error getting sensor", "error", err)
return Sensor{}, err
}
return s.repo.ReadSensor(sensor.SensorID)
} }
func (s *Service) GetValues(sensorID string, from, to time.Time) ([]SensorData, error) { func (s *Service) GetValues(sensor SensorDataRequest) ([]SensorData, error) {
return s.repo.ReadSensorValues(sensorID, from, to) if err := sensor.Validate(); err != nil {
slog.Error("error validating sensor data request", "error", err)
return []SensorData{}, err
}
from, err := time.Parse(time.RFC3339, *sensor.From)
if err != nil {
slog.Error("error parsing from date", "error", err)
return []SensorData{}, err
}
to, err := time.Parse(time.RFC3339, *sensor.To)
if err != nil {
slog.Error("error parsing to date", "error", err)
return []SensorData{}, err
}
return s.repo.ReadSensorValues(sensor.SensorID, from, to)
} }
func (s *Service) ListSensors() ([]Sensor, error) { func (s *Service) ListSensors() ([]Sensor, error) {

View File

@ -22,6 +22,9 @@ func Start(nats *broker.NATS) *Simulator {
} }
} }
// SimulateSensor simula lo que es un sensor, se llama a ese método como una
// go-rutina separada. Hace uso del SamplingInterval como temporizador para
// el canal ticker.
func (s *Simulator) SimulateSensor(sensor Sensor) { func (s *Simulator) SimulateSensor(sensor Sensor) {
s.mu.Lock() s.mu.Lock()
stopChan := make(chan bool) stopChan := make(chan bool)
@ -60,6 +63,8 @@ func (s *Simulator) SimulateSensor(sensor Sensor) {
} }
} }
// UpdateSensor para la gorutina que haya activa de dicho sensor, y comienza una
// nueva con el intervalo actualizado.
func (s *Simulator) UpdateSensor(sensor Sensor) { func (s *Simulator) UpdateSensor(sensor Sensor) {
s.mu.Lock() s.mu.Lock()
stopChan, exists := s.stopChannels[sensor.SensorID] stopChan, exists := s.stopChannels[sensor.SensorID]
@ -77,6 +82,7 @@ func (s *Simulator) UpdateSensor(sensor Sensor) {
slog.Info("simulator updated for sensor", "sensor_id", sensor.SensorID, "new_interval", sensor.SamplingInterval) slog.Info("simulator updated for sensor", "sensor_id", sensor.SensorID, "new_interval", sensor.SamplingInterval)
} }
// generateData genera datos aleatorios por cada tipo de sensor.
func (s *Simulator) generateData(sensor Sensor) SensorData { func (s *Simulator) generateData(sensor Sensor) SensorData {
now := time.Now() now := time.Now()
data := SensorData{ data := SensorData{