package stats import ( "context" "database/sql" "math" "regexp" "strings" "testing" "time" "github.com/DATA-DOG/go-sqlmock" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" ) func pressureEnvelope(at time.Time) envelope.FrameEnvelope { return envelope.FrameEnvelope{ Protocol: envelope.ProtocolGB32960, VIN: "LA9GG64L0NBAF4175", EventID: "event-1", EventTimeMS: at.UnixMilli(), SourceEndpoint: "10.0.0.1:9000", Fields: map[string]any{ "gb32960.fuel_cell.max_hydrogen_pressure_mpa": 6.0, "gb32960.fuel_cell.max_hydrogen_temperature_c": 44.0, "gb32960.gd_fc_vehicle_info.hydrogen_mass_kg": 99.9, }, } } func TestHydrogenStreamSampleUsesPressureTemperatureAndCapacity(t *testing.T) { loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 7, 21, 10, 30, 0, 0, loc) sample, reason, ok := HydrogenStreamSampleFromEnvelope(pressureEnvelope(eventTime), loc, eventTime.Add(time.Minute), 2050) if !ok || reason != "" { t.Fatalf("sample rejected: reason=%q", reason) } if sample.Date != "2026-07-21" || sample.EventID != "event-1" || sample.TankCapacityLiters != 2050 { t.Fatalf("unexpected sample: %#v", sample) } if math.Abs(sample.MassKG-9.092) > 0.01 { t.Fatalf("pressure mass = %.3fkg, want about 9.092kg", sample.MassKG) } if sample.MassKG == 99.9 { t.Fatal("reported Guangdong hydrogen mass must be ignored") } _, reason, ok = HydrogenStreamSampleFromEnvelope(pressureEnvelope(eventTime), loc, eventTime.AddDate(0, 0, 1), 2050) if ok || reason != "not_current_date" { t.Fatalf("historical sample accepted: ok=%v reason=%q", ok, reason) } } func TestHydrogenStreamSampleExtractsV35Inputs(t *testing.T) { loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 9, 4, 10, 30, 0, 0, loc) env := pressureEnvelope(eventTime) env.Fields[envelope.FieldTotalMileageKM] = 12005.6 env.Fields[envelope.FieldSOCPercent] = 72.5 env.Fields[envelope.FieldVehicleRunningMode] = 2 env.Fields[envelope.FieldFuelCellWorkMode] = 2 env.Fields["vehicle_status"] = 1 env.Fields["charge_status"] = 3 env.Fields["fuel_cell_voltage_v"] = 480.0 env.Fields["fuel_cell_current_a"] = 100.0 sample, reason, ok := HydrogenStreamSampleFromEnvelope(env, loc, eventTime, 520) if !ok || reason != "" { t.Fatalf("sample rejected: reason=%q", reason) } if !sample.MileageKnown || sample.MileageKM != 12005.6 || !sample.SOCKnown || sample.SOCPercent != 72.5 || !sample.VehicleStateKnown || sample.VehicleState != 1 || !sample.ChargeStateKnown || sample.ChargeState != 3 || !sample.RunningModeKnown || sample.RunningMode != 2 || !sample.FuelCellPowerKnown || sample.FuelCellPowerKW != 48 { t.Fatalf("V3.5 fields=%#v", sample) } } func TestHydrogenStreamMarksInactiveFuelCellIneligible(t *testing.T) { loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 8, 15, 13, 31, 7, 0, loc) env := pressureEnvelope(eventTime) env.Fields[envelope.FieldFuelCellWorkMode] = 0 sample, reason, ok := HydrogenStreamSampleFromEnvelope(env, loc, eventTime, 520) if !ok || reason != "" || !sample.FuelCellStateKnown || sample.FuelCellActive || sample.ConsumptionEligible { t.Fatalf("inactive sample classification=%#v reason=%q ok=%v", sample, reason, ok) } } func TestHydrogenStreamFallsBackToFuelCellCurrent(t *testing.T) { loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 8, 15, 13, 31, 7, 0, loc) env := pressureEnvelope(eventTime) delete(env.Fields, envelope.FieldFuelCellWorkMode) env.Fields["fuel_cell_current_a"] = 0.2 sample, reason, ok := HydrogenStreamSampleFromEnvelope(env, loc, eventTime, 520) if !ok || reason != "" || !sample.FuelCellStateKnown || sample.FuelCellActive || sample.ConsumptionEligible { t.Fatalf("current fallback classification=%#v reason=%q ok=%v", sample, reason, ok) } } func TestHydrogenStreamRejectsMissingCapacity(t *testing.T) { loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 7, 21, 10, 30, 0, 0, loc) _, reason, ok := HydrogenStreamSampleFromEnvelope(pressureEnvelope(eventTime), loc, eventTime, 0) if ok || reason != "missing_tank_capacity" { t.Fatalf("missing capacity accepted: ok=%v reason=%q", ok, reason) } } func TestHydrogenStreamRejectsInvalidPressureTemperaturePlaceholders(t *testing.T) { loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 8, 12, 8, 11, 16, 0, loc) for _, fields := range []struct{ pressure, temperature float64 }{{0, 30}, {22.8, -40}} { env := pressureEnvelope(eventTime) env.Fields["gb32960.fuel_cell.max_hydrogen_pressure_mpa"] = fields.pressure env.Fields["gb32960.fuel_cell.max_hydrogen_temperature_c"] = fields.temperature _, reason, ok := HydrogenStreamSampleFromEnvelope(env, loc, eventTime, 520) if ok || reason != "invalid_pressure_temperature" { t.Fatalf("placeholder pressure=%v temperature=%v accepted: ok=%v reason=%q", fields.pressure, fields.temperature, ok, reason) } } } func TestNISTHydrogenDensityValidationPoint(t *testing.T) { density, ok := HydrogenDensityKGPerM3(10, 26.85) if !ok || math.Abs(density-7.625) > 0.01 { t.Fatalf("density=%.6f ok=%v", density, ok) } } func hydrogenSegmentTestSample(at time.Time, mass float64, source string, active bool) HydrogenStreamSample { return HydrogenStreamSample{ VIN: "LA9GG64L0NBAF4175", Date: at.Format("2006-01-02"), SourceEndpoint: source, EventID: at.Format(time.RFC3339Nano), EventTime: at, MassKG: mass, TankCapacityLiters: 520, PressureMPa: mass, TemperatureC: 30, NoiseKG: 0.05, RefuelThresholdKG: 1, FuelCellActive: active, FuelCellStateKnown: true, ConsumptionEligible: active, } } func TestHydrogenSegmentStreamUsesEndpointMediansAcrossSources(t *testing.T) { base := time.Date(2026, 8, 19, 8, 0, 0, 0, time.FixedZone("Asia/Shanghai", 8*3600)) state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 10, "source-a", true), nil) for index := 1; index < 10; index++ { mass := 10.0 if index >= 5 { mass = 9.6 } source := "source-a" if index >= 6 { source = "source-b" } if !state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), mass, source, true)) { t.Fatalf("sample %d rejected", index) } } quality, reason := state.quality() if state.projectedConsumptionKG() != 0.4 || quality != "OK" || reason != "" || state.SampleCount != 10 { t.Fatalf("state=%#v consumption=%.3f quality=%s reason=%q", state, state.projectedConsumptionKG(), quality, reason) } if state.SourceEndpoint != "source-b" || len(state.Segment.First) != 5 || len(state.Segment.Tail) != 5 { t.Fatalf("endpoint/window state=%#v", state) } } func TestHydrogenSegmentStreamFinalizesOnInactiveFrame(t *testing.T) { base := time.Now() state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 10, "source-a", true), nil) for index := 1; index < 10; index++ { mass := 10.0 if index >= 5 { mass = 9.6 } state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), mass, "source-a", true)) } state.add(hydrogenSegmentTestSample(base.Add(10*time.Second), 9.5, "source-a", false)) if math.Abs(state.FinalizedConsumptionKG-0.4) > 0.0001 || state.Segment.Count != 0 || state.QualifiedSegmentCount != 1 { t.Fatalf("state=%#v", state) } } func TestHydrogenSegmentStreamIgnoresOutOfOrderAndKeepsBoundedWindow(t *testing.T) { base := time.Now() state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 20, "source-a", true), nil) if state.add(hydrogenSegmentTestSample(base, 19, "source-b", true)) { t.Fatal("same-time duplicate was accepted") } for index := 1; index <= 10000; index++ { state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), 20-float64(index)/20000, "source-b", true)) } if len(state.Segment.First) != 5 || len(state.Segment.Tail) != 5 || state.SampleCount != 10001 { t.Fatalf("first=%d tail=%d samples=%d", len(state.Segment.First), len(state.Segment.Tail), state.SampleCount) } } func TestHydrogenSegmentStreamInactiveDropDoesNotConsume(t *testing.T) { base := time.Now() state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 10, "source-a", false), nil) for index := 1; index < 20; index++ { state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), 10-float64(index)/10, "source-a", false)) } quality, _ := state.quality() if state.projectedConsumptionKG() != 0 || state.RemainingMassKG != 10 || quality != "NO_DATA" { t.Fatalf("state=%#v quality=%s", state, quality) } } func TestHydrogenSegmentStreamSplitsAndAccumulatesRefuelCycles(t *testing.T) { base := time.Now() state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 10, "source-a", true), nil) for index := 1; index < 10; index++ { mass := 10.0 if index >= 5 { mass = 9.6 } state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), mass, "source-a", true)) } state.add(hydrogenSegmentTestSample(base.Add(10*time.Second), 12, "source-b", true)) for index := 11; index < 20; index++ { mass := 12.0 if index >= 15 { mass = 11.7 } state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), mass, "source-b", true)) } if got := state.projectedConsumptionKG(); math.Abs(got-0.7) > 0.0001 { t.Fatalf("consumption=%.3f state=%#v", got, state) } if state.QualifiedSegmentCount != 1 || state.RefuelCount != 1 || state.Segment.Count != 10 { t.Fatalf("state=%#v", state) } } func TestHydrogenSegmentStreamFiltersAbnormalDrop(t *testing.T) { base := time.Now() state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 30, "source-a", true), nil) state.add(hydrogenSegmentTestSample(base.Add(time.Second), 5, "source-a", true)) for index := 2; index < 12; index++ { state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), 5, "source-a", true)) } quality, _ := state.quality() if state.AbnormalDropCount != 1 || state.projectedConsumptionKG() != 0 || quality != "SUSPECT" { t.Fatalf("state=%#v quality=%s", state, quality) } } func TestHydrogenSegmentStreamSeedsDayEndBaselineWithoutRecounting(t *testing.T) { base := time.Now() baseline := &hydrogenStreamBaseline{ SourceEndpoint: "source-a", ConsumptionKG: 3.2, FirstMassKG: 14, LastMassKG: 10, SampleCount: 100, RefuelCount: 1, QualityStatus: "OK", } state := newHydrogenSegmentStreamState(hydrogenSegmentTestSample(base, 10, "source-b", true), baseline) for index := 1; index < 10; index++ { mass := 10.0 if index >= 5 { mass = 9.8 } state.add(hydrogenSegmentTestSample(base.Add(time.Duration(index)*time.Second), mass, "source-b", true)) } if got := state.projectedConsumptionKG(); math.Abs(got-3.4) > 0.0001 { t.Fatalf("consumption=%.3f state=%#v", got, state) } if state.FirstMassKG != 14 || state.SampleCount != 110 || state.RefuelCount != 1 { t.Fatalf("state=%#v", state) } } func TestAppendHydrogenStreamPersistsPressureEvidenceAtomically(t *testing.T) { db, mock, err := sqlmock.New() if err != nil { t.Fatal(err) } defer db.Close() loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Date(2026, 7, 21, 10, 30, 0, 0, loc) mock.ExpectBegin() mock.ExpectQuery(regexp.QuoteMeta("SELECT state_json")). WithArgs("LA9GG64L0NBAF4175", "2026-07-21"). WillReturnError(sql.ErrNoRows) mock.ExpectQuery(regexp.QuoteMeta("SELECT source_endpoint,consumption_kg,first_mass_kg,last_mass_kg,")). WithArgs("LA9GG64L0NBAF4175", "2026-07-21"). WillReturnError(sql.ErrNoRows) mock.ExpectExec(regexp.QuoteMeta("INSERT INTO vehicle_open_hydrogen_segment_stream_state")). WithArgs("LA9GG64L0NBAF4175", "2026-07-21", "10.0.0.1:9000", sqlmock.AnyArg(), sqlmock.AnyArg(), int64(1), int64(0), int64(0), int64(0), int64(0), sqlmock.AnyArg(), eventTime, "event-1", sqlmock.AnyArg(), "NO_DATA", "有效车载氢量样本不足2条"). WillReturnResult(sqlmock.NewResult(1, 1)) mock.ExpectExec(regexp.QuoteMeta("INSERT INTO vehicle_open_daily_energy")). WithArgs("LA9GG64L0NBAF4175", "2026-07-21", "10.0.0.1:9000", 0.0, sqlmock.AnyArg(), sqlmock.AnyArg(), int64(1), int64(0), "NO_DATA", "有效车载氢量样本不足2条"). WillReturnResult(sqlmock.NewResult(1, 1)) mock.ExpectCommit() result, err := AppendHydrogenStream(context.Background(), db, pressureEnvelope(eventTime), loc, eventTime, 2050) if err != nil { t.Fatal(err) } if result.Found != 1 || result.Written != 1 { t.Fatalf("unexpected result: %#v", result) } if err := mock.ExpectationsWereMet(); err != nil { t.Fatal(err) } } func TestWriterAppendWithResultUsesInMemoryCapacity(t *testing.T) { db, mock, err := sqlmock.New() if err != nil { t.Fatal(err) } defer db.Close() loc := time.FixedZone("Asia/Shanghai", 8*3600) eventTime := time.Now().In(loc).Truncate(time.Millisecond) mock.ExpectBegin() mock.ExpectQuery(regexp.QuoteMeta("SELECT state_json")).WillReturnError(sql.ErrNoRows) mock.ExpectExec(regexp.QuoteMeta("INSERT INTO vehicle_open_hydrogen_segment_stream_state")).WillReturnResult(sqlmock.NewResult(1, 1)) mock.ExpectExec(regexp.QuoteMeta("INSERT INTO vehicle_open_daily_energy")).WillReturnResult(sqlmock.NewResult(1, 1)) mock.ExpectCommit() writer := NewWriter(db, loc) writer.hydrogenTankCapacities["LA9GG64L0NBAF4175"] = 2050 writer.hydrogenEnergyParameters["LA9GG64L0NBAF4175"] = HydrogenRealtimeParameters{BatteryCapacityKWh: 80, HydrogenEnergyKWhKG: 16} result, err := writer.AppendWithResult(context.Background(), pressureEnvelope(eventTime)) if err != nil { t.Fatalf("AppendWithResult() error = %v", err) } if result.HydrogenSamplesFound != 1 || result.HydrogenSamplesWritten != 1 { t.Fatalf("hydrogen counters lost: %+v", result) } if err := mock.ExpectationsWereMet(); err != nil { t.Fatal(err) } } func TestHydrogenStreamSQLUsesVINDateStateAndSegmentMethod(t *testing.T) { for _, want := range []string{ "PRIMARY KEY (vin, stat_date)", "state_json JSON NOT NULL", HydrogenV35RealtimeAlgorithmVersion, } { if !strings.Contains(HydrogenSegmentStreamStateTableSQL, want) { t.Fatalf("segment stream schema missing %q", want) } } if !strings.Contains(selectHydrogenSegmentStreamStateSQL, "FOR UPDATE") { t.Fatal("segment state must be locked before update") } }