Все три Go-сервиса (ingest, device-control, rule-engine) отдают /metrics в формате Prometheus: счётчики MQTT/RabbitMQ/ClickHouse операций, гистограммы длительности батч-флашей и диспатча правил, переходы устройств online/offline. Добавлены сервисы prometheus и grafana в docker-compose с провижининг конфигом (datasource + два готовых дашборда: "Ingest & Telemetry" и "Automation & Devices").
187 lines
6.9 KiB
Go
187 lines
6.9 KiB
Go
package engine
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"io"
|
|
"log/slog"
|
|
"testing"
|
|
|
|
"github.com/prometheus/client_golang/prometheus/testutil"
|
|
|
|
"git.cactoz.su/cacto/home_automatization/services/rule-engine-service/internal/devicecontrolclient"
|
|
"git.cactoz.su/cacto/home_automatization/services/rule-engine-service/internal/metrics"
|
|
"git.cactoz.su/cacto/home_automatization/services/rule-engine-service/internal/rules"
|
|
)
|
|
|
|
type fakeRuleSource struct {
|
|
rules []rules.Rule
|
|
}
|
|
|
|
func (f *fakeRuleSource) MatchingRules(_, _ string) []rules.Rule {
|
|
return f.rules
|
|
}
|
|
|
|
type dispatchCall struct {
|
|
deviceID string
|
|
actionType string
|
|
params map[string]any
|
|
}
|
|
|
|
type fakeDispatcher struct {
|
|
calls []dispatchCall
|
|
result devicecontrolclient.Result
|
|
err error
|
|
}
|
|
|
|
func (f *fakeDispatcher) Dispatch(_ context.Context, deviceID, actionType string, params map[string]any) (devicecontrolclient.Result, error) {
|
|
f.calls = append(f.calls, dispatchCall{deviceID: deviceID, actionType: actionType, params: params})
|
|
return f.result, f.err
|
|
}
|
|
|
|
type triggeredEvent struct {
|
|
ruleID int64
|
|
deviceID string
|
|
actionType string
|
|
success bool
|
|
errMsg string
|
|
}
|
|
|
|
type fakePublisher struct {
|
|
events []triggeredEvent
|
|
}
|
|
|
|
func (f *fakePublisher) PublishRuleTriggered(_ context.Context, ruleID, _ int64, deviceID, actionType string, success bool, errMsg string) error {
|
|
f.events = append(f.events, triggeredEvent{ruleID: ruleID, deviceID: deviceID, actionType: actionType, success: success, errMsg: errMsg})
|
|
return nil
|
|
}
|
|
|
|
func discardLogger() *slog.Logger {
|
|
return slog.New(slog.NewTextHandler(io.Discard, nil))
|
|
}
|
|
|
|
func TestHandleReading_NoMatchingRules(t *testing.T) {
|
|
dispatcher := &fakeDispatcher{}
|
|
publisher := &fakePublisher{}
|
|
e := New(&fakeRuleSource{}, dispatcher, publisher, discardLogger())
|
|
|
|
if err := e.HandleReading(context.Background(), Reading{DeviceID: "sensor-1", SensorType: "temperature", Value: 30}); err != nil {
|
|
t.Fatalf("unexpected error: %v", err)
|
|
}
|
|
if len(dispatcher.calls) != 0 || len(publisher.events) != 0 {
|
|
t.Fatalf("expected no dispatch/publish, got %+v / %+v", dispatcher.calls, publisher.events)
|
|
}
|
|
}
|
|
|
|
func TestHandleReading_ConditionNotMet(t *testing.T) {
|
|
source := &fakeRuleSource{rules: []rules.Rule{
|
|
{ID: 1, ConditionOperator: ">", ConditionValue: 28, TargetDeviceID: "fan-1", ActionType: "turn_on"},
|
|
}}
|
|
dispatcher := &fakeDispatcher{}
|
|
e := New(source, dispatcher, &fakePublisher{}, discardLogger())
|
|
|
|
if err := e.HandleReading(context.Background(), Reading{Value: 20}); err != nil {
|
|
t.Fatalf("unexpected error: %v", err)
|
|
}
|
|
if len(dispatcher.calls) != 0 {
|
|
t.Fatalf("expected no dispatch when condition not met, got %+v", dispatcher.calls)
|
|
}
|
|
}
|
|
|
|
func TestHandleReading_ConditionMet_DispatchesAndPublishes(t *testing.T) {
|
|
source := &fakeRuleSource{rules: []rules.Rule{
|
|
{ID: 1, ZoneID: 7, ConditionOperator: ">", ConditionValue: 28, TargetDeviceID: "fan-1", ActionType: "turn_on"},
|
|
}}
|
|
dispatcher := &fakeDispatcher{result: devicecontrolclient.Result{Success: true}}
|
|
publisher := &fakePublisher{}
|
|
e := New(source, dispatcher, publisher, discardLogger())
|
|
|
|
if err := e.HandleReading(context.Background(), Reading{Value: 30}); err != nil {
|
|
t.Fatalf("unexpected error: %v", err)
|
|
}
|
|
|
|
if len(dispatcher.calls) != 1 || dispatcher.calls[0].deviceID != "fan-1" || dispatcher.calls[0].actionType != "turn_on" {
|
|
t.Fatalf("got dispatch calls %+v, want one turn_on for fan-1", dispatcher.calls)
|
|
}
|
|
if len(publisher.events) != 1 || !publisher.events[0].success {
|
|
t.Fatalf("got publish events %+v, want one success event", publisher.events)
|
|
}
|
|
}
|
|
|
|
func TestHandleReading_BusinessRejection_NotRetried(t *testing.T) {
|
|
source := &fakeRuleSource{rules: []rules.Rule{
|
|
{ID: 1, ConditionOperator: ">", ConditionValue: 28, TargetDeviceID: "fan-1", ActionType: "turn_on"},
|
|
}}
|
|
dispatcher := &fakeDispatcher{result: devicecontrolclient.Result{Success: false, Error: "device offline"}}
|
|
publisher := &fakePublisher{}
|
|
e := New(source, dispatcher, publisher, discardLogger())
|
|
|
|
err := e.HandleReading(context.Background(), Reading{Value: 30})
|
|
if err != nil {
|
|
t.Fatalf("business rejection should not be reported as a retryable error, got %v", err)
|
|
}
|
|
if len(publisher.events) != 1 || publisher.events[0].success || publisher.events[0].errMsg != "device offline" {
|
|
t.Fatalf("got publish events %+v, want one failed event with device offline", publisher.events)
|
|
}
|
|
}
|
|
|
|
func TestHandleReading_TransportError_IsRetryable(t *testing.T) {
|
|
source := &fakeRuleSource{rules: []rules.Rule{
|
|
{ID: 1, ConditionOperator: ">", ConditionValue: 28, TargetDeviceID: "fan-1", ActionType: "turn_on"},
|
|
}}
|
|
dispatcher := &fakeDispatcher{err: errors.New("connection refused")}
|
|
publisher := &fakePublisher{}
|
|
e := New(source, dispatcher, publisher, discardLogger())
|
|
|
|
err := e.HandleReading(context.Background(), Reading{Value: 30})
|
|
if err == nil {
|
|
t.Fatal("expected a retryable error when dispatch fails at the transport level")
|
|
}
|
|
if len(publisher.events) != 1 || publisher.events[0].success {
|
|
t.Fatalf("got publish events %+v, want one failed event recorded even on transport error", publisher.events)
|
|
}
|
|
}
|
|
|
|
func TestHandleReading_InvalidOperatorSkipsRuleButContinues(t *testing.T) {
|
|
source := &fakeRuleSource{rules: []rules.Rule{
|
|
{ID: 1, ConditionOperator: "~=", ConditionValue: 28, TargetDeviceID: "bad-rule"},
|
|
{ID: 2, ConditionOperator: ">", ConditionValue: 28, TargetDeviceID: "fan-1", ActionType: "turn_on"},
|
|
}}
|
|
dispatcher := &fakeDispatcher{result: devicecontrolclient.Result{Success: true}}
|
|
e := New(source, dispatcher, &fakePublisher{}, discardLogger())
|
|
|
|
if err := e.HandleReading(context.Background(), Reading{Value: 30}); err != nil {
|
|
t.Fatalf("unexpected error: %v", err)
|
|
}
|
|
if len(dispatcher.calls) != 1 || dispatcher.calls[0].deviceID != "fan-1" {
|
|
t.Fatalf("got dispatch calls %+v, want only the valid rule dispatched", dispatcher.calls)
|
|
}
|
|
}
|
|
|
|
func TestHandleReading_RecordsMetrics(t *testing.T) {
|
|
source := &fakeRuleSource{rules: []rules.Rule{
|
|
{ID: 1, ConditionOperator: ">", ConditionValue: 28, TargetDeviceID: "fan-1", ActionType: "turn_on"},
|
|
}}
|
|
dispatcher := &fakeDispatcher{result: devicecontrolclient.Result{Success: true}}
|
|
e := New(source, dispatcher, &fakePublisher{}, discardLogger())
|
|
|
|
triggeredCounter := metrics.RulesTriggered.WithLabelValues("turn_on", "success")
|
|
beforeConsumed := testutil.ToFloat64(metrics.ReadingsConsumed)
|
|
beforeMatched := testutil.ToFloat64(metrics.RulesMatched)
|
|
beforeTriggered := testutil.ToFloat64(triggeredCounter)
|
|
|
|
if err := e.HandleReading(context.Background(), Reading{Value: 30}); err != nil {
|
|
t.Fatalf("unexpected error: %v", err)
|
|
}
|
|
|
|
if got := testutil.ToFloat64(metrics.ReadingsConsumed); got != beforeConsumed+1 {
|
|
t.Fatalf("got readings_consumed %v, want %v", got, beforeConsumed+1)
|
|
}
|
|
if got := testutil.ToFloat64(metrics.RulesMatched); got != beforeMatched+1 {
|
|
t.Fatalf("got rules_matched %v, want %v", got, beforeMatched+1)
|
|
}
|
|
if got := testutil.ToFloat64(triggeredCounter); got != beforeTriggered+1 {
|
|
t.Fatalf("got rules_triggered{turn_on,success} %v, want %v", got, beforeTriggered+1)
|
|
}
|
|
}
|