855 lines
26 KiB
Go
855 lines
26 KiB
Go
package realtime
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/alicebob/miniredis/v2"
|
|
"github.com/redis/go-redis/v9"
|
|
|
|
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
|
|
)
|
|
|
|
func TestRepositoryUpdatesMergedAndProtocolSnapshots(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Fields: map[string]any{
|
|
envelope.FieldLongitude: 121.1,
|
|
envelope.FieldTotalMileageKM: 10.5,
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 900,
|
|
ReceivedAtMS: 1200,
|
|
Fields: map[string]any{
|
|
envelope.FieldLongitude: 120.0,
|
|
envelope.FieldLatitude: 30.5,
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
merged, err := repo.GetMerged(ctx, "VIN001")
|
|
if err != nil {
|
|
t.Fatalf("GetMerged() error = %v", err)
|
|
}
|
|
if merged.Fields[envelope.FieldLongitude] != 121.1 {
|
|
t.Fatalf("older longitude overwrote newer value: %#v", merged.Fields)
|
|
}
|
|
if merged.Fields[envelope.FieldLatitude] != 30.5 {
|
|
t.Fatalf("new latitude missing: %#v", merged.Fields)
|
|
}
|
|
if got := protocolNames(merged.Protocols); strings.Join(got, ",") != "GB32960,JT808" {
|
|
t.Fatalf("merged protocols = %#v", got)
|
|
}
|
|
protocol, err := repo.GetProtocol(ctx, "VIN001", envelope.ProtocolJT808)
|
|
if err != nil {
|
|
t.Fatalf("GetProtocol() error = %v", err)
|
|
}
|
|
if protocol.Fields[envelope.FieldTotalMileageKM] != 10.5 {
|
|
t.Fatalf("protocol snapshot = %#v", protocol)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryStoresFullParsedOnlyInRealtimeRawAndMergesGB32960Units(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"header": map[string]any{"command": "0x02"},
|
|
"data_units": []any{
|
|
map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}},
|
|
},
|
|
},
|
|
Fields: map[string]any{envelope.FieldSOCPercent: 88.0},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1100,
|
|
ReceivedAtMS: 1200,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{"type": "0x05", "name": "position", "value": map[string]any{"longitude": 121.1, "latitude": 30.2}},
|
|
},
|
|
},
|
|
Fields: map[string]any{envelope.FieldLongitude: 121.1, envelope.FieldLatitude: 30.2},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
protocol, err := repo.GetProtocol(ctx, "VIN001", envelope.ProtocolGB32960)
|
|
if err != nil {
|
|
t.Fatalf("GetProtocol() error = %v", err)
|
|
}
|
|
realtimeRaw, err := repo.GetRealtimeRaw(ctx, "VIN001", envelope.ProtocolGB32960)
|
|
if err != nil {
|
|
t.Fatalf("GetRealtimeRaw() error = %v", err)
|
|
}
|
|
protocolJSON, err := json.Marshal(protocol)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if strings.Contains(string(protocolJSON), `"parsed"`) {
|
|
t.Fatalf("protocol snapshot should not duplicate full parsed data: %s", string(protocolJSON))
|
|
}
|
|
if len(realtimeRaw["data_units"].([]any)) != 2 {
|
|
t.Fatalf("realtime-raw did not keep merged parsed fields: %#v", realtimeRaw)
|
|
}
|
|
merged, err := repo.GetMerged(ctx, "VIN001")
|
|
if err != nil {
|
|
t.Fatalf("GetMerged() error = %v", err)
|
|
}
|
|
mergedJSON, err := json.Marshal(merged)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if strings.Contains(string(mergedJSON), "protocol_data") {
|
|
t.Fatalf("merged snapshot should not duplicate full protocol parsed data: %s", string(mergedJSON))
|
|
}
|
|
}
|
|
|
|
func TestRepositoryMergesNestedGB32960MotorSlicesBySerialNo(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{
|
|
"type": "0x02",
|
|
"name": "drive_motor",
|
|
"value": map[string]any{
|
|
"motors": []any{
|
|
map[string]any{"serial_no": 1, "speed_rpm": 5000},
|
|
map[string]any{"serial_no": 2, "speed_rpm": 5100},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1100,
|
|
ReceivedAtMS: 1200,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{
|
|
"type": "0x02",
|
|
"name": "drive_motor",
|
|
"value": map[string]any{
|
|
"motors": []any{
|
|
map[string]any{"serial_no": 1, "speed_rpm": 5406},
|
|
map[string]any{"serial_no": 2, "speed_rpm": 5377},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
realtimeRaw, err := repo.GetRealtimeRaw(ctx, "VIN001", envelope.ProtocolGB32960)
|
|
if err != nil {
|
|
t.Fatalf("GetRealtimeRaw() error = %v", err)
|
|
}
|
|
units := realtimeRaw["data_units"].([]any)
|
|
driveMotor := units[0].(map[string]any)
|
|
value := driveMotor["value"].(map[string]any)
|
|
motors := value["motors"].([]any)
|
|
if len(motors) != 2 {
|
|
t.Fatalf("motors len = %d, want 2: %#v", len(motors), motors)
|
|
}
|
|
first := motors[0].(map[string]any)
|
|
if first["speed_rpm"] != float64(5406) {
|
|
t.Fatalf("first motor was not updated: %#v", motors)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryCollapsesGB32960VendorStackFragmentsIntoSingleRealtimeSummary(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
first := map[string]any{
|
|
"engine_work_state": 1,
|
|
"stack_water_outlet_temp_c": 65,
|
|
"hydrogen_inlet_pressure_kpa": 130,
|
|
"cell_count": 432,
|
|
"avg_cell_voltage_v": 0.78,
|
|
"frame_cell_start": 201,
|
|
"frame_cell_count": 200,
|
|
"frame_max_cell_voltage_v": 1.0,
|
|
"frame_min_cell_voltage_v": 0.78,
|
|
}
|
|
second := map[string]any{
|
|
"engine_work_state": 1,
|
|
"stack_water_outlet_temp_c": 63,
|
|
"hydrogen_inlet_pressure_kpa": 130,
|
|
"cell_count": 432,
|
|
"avg_cell_voltage_v": 0.80,
|
|
"frame_cell_start": 401,
|
|
"frame_cell_count": 32,
|
|
"frame_max_cell_voltage_v": 1.0,
|
|
"frame_min_cell_voltage_v": 1.0,
|
|
}
|
|
for i, summary := range []map[string]any{first, second} {
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: int64(1000 + i),
|
|
ReceivedAtMS: int64(1100 + i),
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{
|
|
"type": "0x30",
|
|
"name": "gd_fc_stack",
|
|
"value": map[string]any{
|
|
"stack_count": 1,
|
|
"summaries": []any{summary},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
}
|
|
|
|
realtimeRaw, err := repo.GetRealtimeRaw(ctx, "VIN001", envelope.ProtocolGB32960)
|
|
if err != nil {
|
|
t.Fatalf("GetRealtimeRaw() error = %v", err)
|
|
}
|
|
unit := realtimeRaw["data_units"].([]any)[0].(map[string]any)
|
|
value := unit["value"].(map[string]any)
|
|
summaries := value["summaries"].([]any)
|
|
if len(summaries) != 1 {
|
|
t.Fatalf("stack summaries should be a current stack state, not appended fragments: %#v", summaries)
|
|
}
|
|
summary := summaries[0].(map[string]any)
|
|
for _, fragmentOnly := range []string{"frame_cell_start", "frame_cell_count", "frame_max_cell_voltage_v", "frame_min_cell_voltage_v"} {
|
|
if _, ok := summary[fragmentOnly]; ok {
|
|
t.Fatalf("stack summary should not expose fragment-only field %s: %#v", fragmentOnly, summary)
|
|
}
|
|
}
|
|
if summary["stack_water_outlet_temp_c"] != float64(63) {
|
|
t.Fatalf("summary should keep latest stack-level values: %#v", summary)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryReplacesGB32960VendorAuxiliarySubsystemsWithLatestState(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
for i, voltage := range []float64{372, 348} {
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: int64(1000 + i),
|
|
ReceivedAtMS: int64(1100 + i),
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{
|
|
"type": "0x31",
|
|
"name": "gd_fc_auxiliary",
|
|
"value": map[string]any{
|
|
"subsystem_count": 1,
|
|
"subsystems": []any{map[string]any{
|
|
"air_compressor_motor_voltage_v": voltage,
|
|
"water_pump_voltage_v": voltage,
|
|
}},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
}
|
|
|
|
realtimeRaw, err := repo.GetRealtimeRaw(ctx, "VIN001", envelope.ProtocolGB32960)
|
|
if err != nil {
|
|
t.Fatalf("GetRealtimeRaw() error = %v", err)
|
|
}
|
|
unit := realtimeRaw["data_units"].([]any)[0].(map[string]any)
|
|
value := unit["value"].(map[string]any)
|
|
subsystems := value["subsystems"].([]any)
|
|
if len(subsystems) != 1 {
|
|
t.Fatalf("auxiliary subsystems should keep current state, not append historical frames: %#v", subsystems)
|
|
}
|
|
subsystem := subsystems[0].(map[string]any)
|
|
if subsystem["air_compressor_motor_voltage_v"] != float64(348) {
|
|
t.Fatalf("auxiliary subsystem should keep latest values: %#v", subsystem)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryWritesRealtimeKVHashes(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}},
|
|
map[string]any{"type": "0x30", "name": "gd_fc_stack", "value": map[string]any{
|
|
"stack_count": 1,
|
|
"summaries": []any{map[string]any{"stack_water_outlet_temp_c": 63}},
|
|
}},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
dataUnitsKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:values").Result()
|
|
if err != nil {
|
|
t.Fatalf("values kv HGetAll error = %v", err)
|
|
}
|
|
if dataUnitsKV["gb32960.vehicle.soc_percent"] != "88" {
|
|
t.Fatalf("data_units kv = %#v", dataUnitsKV)
|
|
}
|
|
if dataUnitsKV["gb32960.gd_fc_stack.stack_water_outlet_temp_c"] != "63" {
|
|
t.Fatalf("data_units stack kv = %#v", dataUnitsKV)
|
|
}
|
|
typeKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:types").Result()
|
|
if err != nil {
|
|
t.Fatalf("types kv HGetAll error = %v", err)
|
|
}
|
|
if typeKV["gb32960.vehicle.soc_percent"] != "number" || typeKV["gb32960.gd_fc_stack.type"] != "string" {
|
|
t.Fatalf("types kv = %#v", typeKV)
|
|
}
|
|
metaKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:meta").Result()
|
|
if err != nil {
|
|
t.Fatalf("meta kv HGetAll error = %v", err)
|
|
}
|
|
if metaKV["event_id"] == "" || metaKV["event_time_ms"] != "1000" || metaKV["protocol"] != "GB32960" {
|
|
t.Fatalf("meta kv = %#v", metaKV)
|
|
}
|
|
if repo.client.Exists(ctx, "vehicle:rt-kv:GB32960:VIN001:data_units").Val() != 0 {
|
|
t.Fatal("old domain kv key should not be written")
|
|
}
|
|
if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:values").Val(); ttl != -1 {
|
|
t.Fatalf("kv ttl should not expire, got %v", ttl)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryFastUpdateOnlyWritesPermanentKVAndMinuteOnline(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.FastUpdate(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("FastUpdate() error = %v", err)
|
|
}
|
|
|
|
dataUnitsKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:values").Result()
|
|
if err != nil {
|
|
t.Fatalf("values kv HGetAll error = %v", err)
|
|
}
|
|
if dataUnitsKV["gb32960.vehicle.soc_percent"] != "88" {
|
|
t.Fatalf("data_units kv = %#v", dataUnitsKV)
|
|
}
|
|
typeKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:types").Result()
|
|
if err != nil {
|
|
t.Fatalf("types kv HGetAll error = %v", err)
|
|
}
|
|
if typeKV["gb32960.vehicle.soc_percent"] != "number" {
|
|
t.Fatalf("types kv = %#v", typeKV)
|
|
}
|
|
if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:values").Val(); ttl != -1 {
|
|
t.Fatalf("kv ttl should not expire, got %v", ttl)
|
|
}
|
|
if ttl := repo.client.TTL(ctx, "vehicle:online:GB32960:VIN001").Val(); ttl <= 0 || ttl > time.Minute {
|
|
t.Fatalf("online ttl = %v, want within 1 minute", ttl)
|
|
}
|
|
state, err := repo.client.HGetAll(ctx, "vehicle:online-state:GB32960:VIN001").Result()
|
|
if err != nil {
|
|
t.Fatalf("online state HGetAll error = %v", err)
|
|
}
|
|
if state["online"] != "true" || state["last_seen_ms"] != "1100" || state["offline_after_ms"] != "61100" {
|
|
t.Fatalf("online state = %#v", state)
|
|
}
|
|
if repo.client.ZScore(ctx, "vehicle:last_seen", "GB32960:VIN001").Err() != nil {
|
|
t.Fatal("last_seen should index protocol+vin")
|
|
}
|
|
if repo.client.SIsMember(ctx, "vehicle:rt-index:GB32960", "VIN001").Val() != true {
|
|
t.Fatal("rt-index should contain vin by protocol")
|
|
}
|
|
if repo.client.Exists(ctx, "vehicle:latest:VIN001").Val() != 0 {
|
|
t.Fatal("fast update should not write merged snapshot")
|
|
}
|
|
}
|
|
|
|
func TestRepositoryListsOnlineStatusesAndPipelineSummary(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
for _, env := range []envelope.FrameEnvelope{
|
|
{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}}},
|
|
},
|
|
},
|
|
{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN002",
|
|
EventTimeMS: 2000,
|
|
ReceivedAtMS: 2200,
|
|
MessageID: "0x0200",
|
|
Parsed: map[string]any{"location": map[string]any{"longitude": 121.1, "latitude": 30.2}},
|
|
Fields: map[string]any{envelope.FieldLongitude: 121.1, envelope.FieldLatitude: 30.2},
|
|
},
|
|
} {
|
|
if err := repo.FastUpdate(ctx, env); err != nil {
|
|
t.Fatalf("FastUpdate() error = %v", err)
|
|
}
|
|
}
|
|
|
|
items, total, err := repo.ListOnline(ctx, OnlineListQuery{Protocol: envelope.ProtocolJT808, Limit: 10})
|
|
if err != nil {
|
|
t.Fatalf("ListOnline() error = %v", err)
|
|
}
|
|
if total != 1 || len(items) != 1 || items[0].VIN != "VIN002" || !items[0].Online {
|
|
t.Fatalf("jt808 online list total=%d items=%#v", total, items)
|
|
}
|
|
|
|
summary, err := repo.PipelineSummary(ctx)
|
|
if err != nil {
|
|
t.Fatalf("PipelineSummary() error = %v", err)
|
|
}
|
|
got := map[envelope.Protocol]ProtocolPipelineSummary{}
|
|
for _, item := range summary.Protocols {
|
|
got[item.Protocol] = item
|
|
}
|
|
if got[envelope.ProtocolGB32960].IndexedCount != 1 || got[envelope.ProtocolGB32960].OnlineCount != 1 || got[envelope.ProtocolGB32960].LatestSeenMS != 1100 {
|
|
t.Fatalf("gb32960 summary = %#v", got[envelope.ProtocolGB32960])
|
|
}
|
|
if got[envelope.ProtocolJT808].IndexedCount != 1 || got[envelope.ProtocolJT808].OnlineCount != 1 || got[envelope.ProtocolJT808].LatestSeenMS != 2200 {
|
|
t.Fatalf("jt808 summary = %#v", got[envelope.ProtocolJT808])
|
|
}
|
|
}
|
|
|
|
func TestOnlineIndexHandlerReturnsPagedOnlineStatuses(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.FastUpdate(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}}},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("FastUpdate() error = %v", err)
|
|
}
|
|
|
|
handler := NewOnlineIndexHandler(repo)
|
|
request := httptest.NewRequest(http.MethodGet, "/api/realtime/online?protocol=gb32960&limit=5", nil)
|
|
response := httptest.NewRecorder()
|
|
handler.ServeHTTP(response, request)
|
|
|
|
if response.Code != http.StatusOK {
|
|
t.Fatalf("status = %d body=%s", response.Code, response.Body.String())
|
|
}
|
|
body := response.Body.String()
|
|
for _, want := range []string{`"protocol":"GB32960"`, `"vin":"VIN001"`, `"online":true`, `"total":1`} {
|
|
if !strings.Contains(body, want) {
|
|
t.Fatalf("response missing %s: %s", want, body)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestPipelineDebugHandlerReturnsProtocolSummary(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.FastUpdate(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolYutongMQTT,
|
|
VIN: "VIN003",
|
|
EventTimeMS: 3000,
|
|
ReceivedAtMS: 3300,
|
|
Parsed: map[string]any{"data": map[string]any{"TOTAL_MILEAGE": 1}},
|
|
Fields: map[string]any{envelope.FieldTotalMileageKM: 1},
|
|
}); err != nil {
|
|
t.Fatalf("FastUpdate() error = %v", err)
|
|
}
|
|
|
|
handler := NewPipelineDebugHandler(repo)
|
|
request := httptest.NewRequest(http.MethodGet, "/api/debug/pipeline", nil)
|
|
response := httptest.NewRecorder()
|
|
handler.ServeHTTP(response, request)
|
|
|
|
if response.Code != http.StatusOK {
|
|
t.Fatalf("status = %d body=%s", response.Code, response.Body.String())
|
|
}
|
|
body := response.Body.String()
|
|
for _, want := range []string{`"protocol":"YUTONG_MQTT"`, `"indexed_count":1`, `"online_count":1`, `"latest_seen_ms":3300`} {
|
|
if !strings.Contains(body, want) {
|
|
t.Fatalf("response missing %s: %s", want, body)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestRepositoryReplacesJT808LocationWhenUpdatingRealtimeRaw(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
MessageID: "0x0200",
|
|
Parsed: map[string]any{
|
|
"header": map[string]any{
|
|
"message_id": "0x0200",
|
|
"phone_raw": "013307795425",
|
|
"phone_bcd_hex": "013307795425",
|
|
},
|
|
"location": map[string]any{
|
|
"additional": []any{
|
|
map[string]any{"id": "0x01", "parsed": map[string]any{"raw_value_tenth_km": 102412}},
|
|
},
|
|
"total_mileage_km": 10241.2,
|
|
},
|
|
},
|
|
Fields: map[string]any{envelope.FieldTotalMileageKM: 10241.2},
|
|
}); err != nil {
|
|
t.Fatalf("Update() initial error = %v", err)
|
|
}
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 2000,
|
|
ReceivedAtMS: 2100,
|
|
MessageID: "0x0200",
|
|
Parsed: map[string]any{
|
|
"header": map[string]any{
|
|
"message_id": "0x0200",
|
|
"phone": "13307795425",
|
|
},
|
|
"location": map[string]any{
|
|
"total_mileage_km": 10241.3,
|
|
"network_signal_strength": byte(31),
|
|
},
|
|
},
|
|
Fields: map[string]any{envelope.FieldTotalMileageKM: 10241.3},
|
|
}); err != nil {
|
|
t.Fatalf("Update() replacement error = %v", err)
|
|
}
|
|
|
|
realtimeRaw, err := repo.GetRealtimeRaw(ctx, "VIN001", envelope.ProtocolJT808)
|
|
if err != nil {
|
|
t.Fatalf("GetRealtimeRaw() error = %v", err)
|
|
}
|
|
header := realtimeRaw["header"].(map[string]any)
|
|
if _, ok := header["phone_raw"]; ok {
|
|
t.Fatalf("header should be replaced, got %#v", header)
|
|
}
|
|
location := realtimeRaw["location"].(map[string]any)
|
|
if _, ok := location["additional"]; ok {
|
|
t.Fatalf("jt808 location should be replaced, got %#v", location)
|
|
}
|
|
if location["total_mileage_km"] != float64(10241.3) {
|
|
t.Fatalf("mileage = %#v", location["total_mileage_km"])
|
|
}
|
|
}
|
|
|
|
func TestRepositoryOnlineStatus(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
status, err := repo.IsOnline(ctx, "VIN001")
|
|
if err != nil {
|
|
t.Fatalf("IsOnline() error = %v", err)
|
|
}
|
|
if status.Online {
|
|
t.Fatal("empty vin should be offline")
|
|
}
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Fields: map[string]any{envelope.FieldSpeedKMH: 12.3},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
status, err = repo.IsOnline(ctx, "VIN001")
|
|
if err != nil {
|
|
t.Fatalf("IsOnline() error = %v", err)
|
|
}
|
|
if !status.Online || status.LastSeenMS != 1100 {
|
|
t.Fatalf("online status = %#v", status)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryUsesEnvelopeParsedFieldsForRealtimeKV(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
MessageID: "0x0200",
|
|
Parsed: map[string]any{
|
|
"location": map[string]any{"speed_kmh": 99},
|
|
},
|
|
ParsedFields: map[string]any{
|
|
"jt808.location.speed_kmh": "12.3",
|
|
},
|
|
Fields: map[string]any{envelope.FieldSpeedKMH: 12.3},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
values, err := repo.client.HGetAll(ctx, realtimeKVValuesKey(envelope.ProtocolJT808, "VIN001")).Result()
|
|
if err != nil {
|
|
t.Fatalf("HGetAll() error = %v", err)
|
|
}
|
|
if got, want := values["jt808.location.speed_kmh"], "12.3"; got != want {
|
|
t.Fatalf("speed kv = %#v, want %q; all=%#v", got, want, values)
|
|
}
|
|
}
|
|
|
|
func TestRepositorySkipsFramesWithoutVIN(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
Phone: "13307811170",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Fields: map[string]any{
|
|
envelope.FieldSpeedKMH: 23.0,
|
|
envelope.FieldTotalMileageKM: 10003.7,
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
if _, err := repo.GetMerged(ctx, "JT808:13307811170"); err == nil {
|
|
t.Fatal("phone-only frame should not create a merged realtime snapshot")
|
|
}
|
|
if _, err := repo.GetRealtimeRaw(ctx, "JT808:13307811170", envelope.ProtocolJT808); err == nil {
|
|
t.Fatal("phone-only frame should not create realtime-raw")
|
|
}
|
|
status, err := repo.IsOnline(ctx, "JT808:13307811170")
|
|
if err != nil {
|
|
t.Fatalf("IsOnline() error = %v", err)
|
|
}
|
|
if status.Online {
|
|
t.Fatalf("phone-only frame should not mark realtime online: %#v", status)
|
|
}
|
|
}
|
|
|
|
func TestRepositoryDoesNotOverwritePositiveMileageWithZero(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
ctx := context.Background()
|
|
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Fields: map[string]any{
|
|
envelope.FieldTotalMileageKM: 12345.6,
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
if err := repo.Update(ctx, envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 2000,
|
|
ReceivedAtMS: 2100,
|
|
Fields: map[string]any{
|
|
envelope.FieldSpeedKMH: 22.0,
|
|
envelope.FieldTotalMileageKM: 0,
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
merged, err := repo.GetMerged(ctx, "VIN001")
|
|
if err != nil {
|
|
t.Fatalf("GetMerged() error = %v", err)
|
|
}
|
|
if merged.Fields[envelope.FieldTotalMileageKM] != 12345.6 {
|
|
t.Fatalf("zero mileage overwrote positive value: %#v", merged.Fields)
|
|
}
|
|
if merged.Fields[envelope.FieldSpeedKMH] != 22.0 {
|
|
t.Fatalf("new speed should still merge: %#v", merged.Fields)
|
|
}
|
|
}
|
|
|
|
func TestHandlerReturnsMergedSnapshot(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
if err := repo.Update(context.Background(), envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
VIN: "VIN001",
|
|
SourceEndpoint: "115.231.168.135:43625",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Fields: map[string]any{envelope.FieldSpeedKMH: 12.3},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
req := httptest.NewRequest(http.MethodGet, "/api/realtime/vehicles/VIN001", nil)
|
|
rec := httptest.NewRecorder()
|
|
NewHandler(repo).ServeHTTP(rec, req)
|
|
if rec.Code != http.StatusOK {
|
|
t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String())
|
|
}
|
|
if !stringsContains(rec.Body.String(), `"vin":"VIN001"`) {
|
|
t.Fatalf("unexpected body: %s", rec.Body.String())
|
|
}
|
|
if stringsContains(rec.Body.String(), "source_endpoint") {
|
|
t.Fatalf("realtime snapshot should not expose source_endpoint: %s", rec.Body.String())
|
|
}
|
|
}
|
|
|
|
func TestHandlerReturnsNotFoundForPhoneOnlyVehicleKey(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
if err := repo.Update(context.Background(), envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolJT808,
|
|
Phone: "13307811170",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Fields: map[string]any{envelope.FieldSpeedKMH: 12.3},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
req := httptest.NewRequest(http.MethodGet, "/api/realtime/vehicles/JT808:13307811170", nil)
|
|
rec := httptest.NewRecorder()
|
|
NewHandler(repo).ServeHTTP(rec, req)
|
|
if rec.Code != http.StatusNotFound {
|
|
t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String())
|
|
}
|
|
}
|
|
|
|
func TestHandlerReturnsRealtimeRaw(t *testing.T) {
|
|
repo, closeFn := newTestRepository(t)
|
|
defer closeFn()
|
|
if err := repo.Update(context.Background(), envelope.FrameEnvelope{
|
|
Protocol: envelope.ProtocolGB32960,
|
|
VIN: "VIN001",
|
|
EventTimeMS: 1000,
|
|
ReceivedAtMS: 1100,
|
|
Parsed: map[string]any{
|
|
"data_units": []any{
|
|
map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 90.0}},
|
|
},
|
|
},
|
|
}); err != nil {
|
|
t.Fatalf("Update() error = %v", err)
|
|
}
|
|
|
|
req := httptest.NewRequest(http.MethodGet, "/api/realtime/vehicles/VIN001/realtime-raw/GB32960", nil)
|
|
rec := httptest.NewRecorder()
|
|
NewHandler(repo).ServeHTTP(rec, req)
|
|
if rec.Code != http.StatusOK {
|
|
t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String())
|
|
}
|
|
if !stringsContains(rec.Body.String(), `"soc_percent":90`) {
|
|
t.Fatalf("unexpected body: %s", rec.Body.String())
|
|
}
|
|
}
|
|
|
|
func newTestRepository(t *testing.T) (*Repository, func()) {
|
|
t.Helper()
|
|
server, err := miniredis.Run()
|
|
if err != nil {
|
|
t.Fatalf("miniredis.Run() error = %v", err)
|
|
}
|
|
client := redis.NewClient(&redis.Options{Addr: server.Addr()})
|
|
return NewRepository(client, Config{OnlineTTL: time.Minute}), func() {
|
|
_ = client.Close()
|
|
server.Close()
|
|
}
|
|
}
|
|
|
|
func stringsContains(value string, pattern string) bool {
|
|
return strings.Contains(value, pattern)
|
|
}
|
|
|
|
func protocolNames(protocols []envelope.Protocol) []string {
|
|
names := make([]string, 0, len(protocols))
|
|
for _, protocol := range protocols {
|
|
names = append(names, string(protocol))
|
|
}
|
|
return names
|
|
}
|